From bea5a1f1c88111df81a23ba1618c94290e1935cd Mon Sep 17 00:00:00 2001 From: Roger Olsson Date: Sat, 12 Sep 2026 12:31:20 -0700 Subject: [PATCH] fcl-web: harden WebSocket client lifetime and I/O Keep the public client/component model and the existing RFC 6455 codec structure, but replace the pump, transport-I/O, session-lifetime and callback-dispatch internals. Use generation-bound reference-counted sessions, stable pump entries and callback admission gates. Disconnect, reconnect, destruction and synchronized callbacks can no longer retain stale indices or unleased connection pointers. The threaded pump isolates blocking reads with one reader per registered session. Make frame and handshake I/O exact, serialize writes, distinguish EOF, socket errors and intentional interruption, capture platform errors immediately, and make per-session cancellation durable. Add bounded handshake and payload limits plus validation for lengths, masking, control frames, UTF-8 and close codes. Document callback, shutdown and ownership policy in the WebSocket source directory. Add a standalone 37-scenario regression harness, with separate pump-before-client and interrupt-race probes. Verified with FPC 3.3.1 on x86-64 Linux: 37 passed, 0 failed, 0 skipped, 0 hung and 0 crashed; --pump-first and --interrupt-race both passed. Plain TCP and TLS blocked-read termination completed in about 101 ms. --- packages/fcl-web/src/websocket/README.md | 350 ++ packages/fcl-web/src/websocket/fpwebsocket.pp | 1218 +++- .../src/websocket/fpwebsocketclient.pp | 2971 ++++++++-- packages/fcl-web/tests/wsshutdowntest.pas | 4897 +++++++++++++++++ 4 files changed, 8870 insertions(+), 566 deletions(-) create mode 100644 packages/fcl-web/src/websocket/README.md create mode 100644 packages/fcl-web/tests/wsshutdowntest.pas diff --git a/packages/fcl-web/src/websocket/README.md b/packages/fcl-web/src/websocket/README.md new file mode 100644 index 00000000..fd6b7f44 --- /dev/null +++ b/packages/fcl-web/src/websocket/README.md @@ -0,0 +1,350 @@ +# FPC WebSocket client bounded rewrite + +## Status and scope + +The client implementation replaces the WebSocket pump, transport-I/O, connection +lifetime, and callback-dispatch internals. It deliberately retains the public +client/component model and the existing RFC 6455 frame classes instead of +replacing the whole API or codec. + +The changed sources are: + +- `fpwebsocket.pp`: shared transport and frame hardening. +- `fpwebsocketclient.pp`: client sessions, pump scheduling, callback gates, + handshake handling, and teardown. +- `packages/fcl-web/tests/wsshutdowntest.pas`: 37 isolated regression scenarios plus the + `--pump-first` and `--interrupt-race` probes. + +The adjacent `fpcustwsserver.pp`, `fpwebsocketserver.pp`, and `wsupgrader.pp` +sources are unchanged by this work. They were reviewed for compatibility but +are not part of the client rewrite; their remaining lifetime problems are +listed near the end of this document. + +## Design goals + +The rewrite is built around five rules: + +1. A connection object must remain alive until every read and callback using it + has returned. +2. A late operation from an old connection generation must never detach or + notify a reconnected generation. +3. No list index or unleased object pointer may survive an application + callback. +4. Registry, component-state, and transport-write locks must not be held while + application code runs. +5. Stopping the pump may be bounded, but destroying an object must not abandon + threads which can still access it. + +## Internal objects and ownership + +| Object | Purpose | Important ownership rule | +| --- | --- | --- | +| `TWSClientOwnerGate` | Controls admission to callbacks into one client component. | Shared by every session generation; may outlive the component. | +| `TWSClientSession` | Represents exactly one `Connect` generation. | Owns the connection, transport/socket, and handshake request/response. | +| `TWSMessagePumpEntry` | Stable registry identity for one session. | Registry, snapshots, and a reader each hold explicit entry references. | +| `TWSPumpCore` | Worker-visible pump state independent of the component lifetime. | Workers hold core references and briefly acquire pump-owner admission before calling pump code. | +| `TWSClientReaderThread` | Performs blocking reads for one session. | Holds both the entry and core until its final cleanup completes. | + +The current client owns one session reference. A pump entry owns another. +Temporary public operations such as `SendMessage`, `Ping`, `Disconnect`, and +manual `CheckIncoming` acquire their own short session reference. A session is +destroyed only after all of these references have gone away. + +The session generation number is checked during connect, handshake completion, +notification, and teardown. An old reader or callback can therefore retire +only its own connection; it cannot clear the fields of a later reconnect. + +`TWSClientConnection.HandshakeResponse` remains a non-owning compatibility +alias. The session owns the actual response and clears the alias before either +object is destroyed. + +## Connect and handshake lifecycle + +`Connect` reserves a new generation under the client state lock. It then holds +an owner-gate admission while constructing the socket, transport, connection, +and session. The session is published as the client's current generation before +the virtual handshake methods are called. + +Every virtual handshake hook is followed by an exact-generation and open-state +check. This permits a hook to disconnect the client without allowing the old +connect operation to publish itself afterwards. + +Handshake hardening includes: + +- one atomic write-all operation for the HTTP upgrade request; +- an actual HTTP status-line parse and mandatory status `101`; +- case-sensitive `Sec-WebSocket-Accept` validation; +- case-insensitive `Upgrade: websocket` validation; +- token-aware `Connection: Upgrade` validation; +- limits of 8 KiB per response header line, 256 lines, and 64 KiB total; +- the same 8 KiB/256-line/64-KiB aggregate policy for the shared server-side + handshake reader. + +## Pump and reader model + +`TWSThreadMessagePump` has a small driver thread and one blocking reader thread +per registered session. A stalled partial frame on one connection therefore +cannot prevent another connection from receiving messages. + +The driver still invokes the protected virtual +`CheckConnections`/`ReadConnections` pipeline. For the default threaded pump, +`ReadConnections` schedules missing per-session readers; it does not read from +a socket itself. This preserves the extension hook without creating two +consumers for one transport. + +`TWSMessagePumpEntry` replaces list-index-based removal. A callback may remove +an earlier connection or reconnect its own client, but the pump subsequently +removes the exact entry it was processing rather than whichever item has moved +into an old index. + +The compatibility `List` property remains a mirror of registered connections +for descendants which inspect it. It is not the lifetime authority. + +## Disconnect ordering and callbacks + +A terminal session follows this order: + +1. Acquire callback admission before attempting the notification claim. This + prevents destruction from closing the gate in the gap after a successful + claim. +2. Atomically claim the session's single disconnect notification. +3. Mark that exact session closing. +4. Clear the client's current fields under its state lock if this is still the + current generation. +5. Remove the exact pump entry. Removal publishes durable read cancellation. +6. Repeat cancellation idempotently, invoke `OnDisconnect` as the tail + operation, and release the client's session ownership. + +Callbacks are admitted through `TWSClientOwnerGate` before the client pointer +is used. Client destruction closes admission and waits for already-admitted +callbacks. When destruction runs on the main thread, the wait services +`CheckSynchronize`, allowing a reader blocked in `TThread.Synchronize` to +finish. + +Supported callback operations: + +| Operation from a callback | Policy | +| --- | --- | +| Disconnect the same or another managed client | Supported. | +| Reconnect from `OnDisconnect` | Supported; it creates a new generation. | +| Call `Terminate` from a reader callback | Supported; it requests stop and does not join itself. | +| Call `Terminate` from a synchronized method | Bounded; unfinished objects are retained and reaped later. | +| Free a client or pump from another thread while its callback runs | Supported; destruction waits for quiescence. | +| Directly free the callback's own client or pump | Rejected with `EWebSocketClient`; returning through a destroyed Pascal object cannot be made safe. | + +Exceptions from `OnDisconnect` and `OnError` cannot skip internal releases. +`OnError` itself is isolated so an error reporter cannot terminate a reader and +discard a pending disconnect. + +Callbacks are not automatically synchronized or serialized onto one thread: + +- handshake callbacks and `OnConnect` run in the thread calling `Connect`; +- `OnMessageReceived` and `OnControl` run in the session reader, or in the + caller when manual `CheckIncoming` is used; +- `OnDisconnect` runs on whichever path wins the once-only notification claim. + +Handshake and `OnConnect` exceptions propagate from and fail `Connect`. With +the default threaded pump, a message/control exception retires the session and +is passed to `OnError`; manual `CheckIncoming` notifies the disconnect and +re-raises the original exception to its caller. + +## Pump termination + +`Terminate` is a two-phase operation: + +1. Stop the pump generation and ask the driver/readers to leave normally. +2. Wait for `2 * Interval + 10 ms`, clamped to 100–1000 ms. +3. If a reader is still blocked in an exact read, repeatedly interrupt only + those active readers and wait for the same bounded interval again. + +A reader interrupted in the middle of a frame is terminal because its parser +cannot resume from a partially consumed frame. Healthy readers normally leave +without interruption, so their entries remain registered and can be served by +a later `Execute`. + +When `Terminate` is running inside a method that the reader synchronized to the +calling thread, it returns after the bounded wait rather than deadlocking. The +thread, core, entry, and session remain leased until the synchronized callback +unwinds. A later `Terminate`, `Execute`, or destructor reaps that generation. + +The destructor is intentionally stronger than `Terminate`: it waits until all +workers have stopped before freeing the component. Abandoning a live worker +would convert a bounded shutdown into a use-after-free. + +## Transport and frame I/O + +`IWSTransport` is byte-for-byte unchanged. The concrete transport internals now +provide: + +- exact reads for frame headers, extended lengths, mask keys, and payloads; +- write-all behavior for short socket writes; +- one transport write lock, preventing frames and handshake bytes from + interleaving; +- immediate capture of the platform socket error before another operation can + overwrite it; +- transient `InterruptRead` for pump shutdown; +- durable per-session `CancelReads`, closing the gap between consecutive exact + read chunks; +- one descriptor close, owned by `TSocketStream`; shutdown wakes blocked I/O + without a later reused-handle double close. + +The relevant exception classes are: + +- `EWSConnectionClosed`: orderly EOF; terminal but not reported as `OnError`; +- `EWSReadError`: terminal socket/read error, reported once; +- `EWSWriteError`: terminal write failure; +- `EWSReadInterrupted`: intentional cancellation/interruption; +- `EWSProtocolError`: RFC violation plus the close status to send. + +On Windows, a reset during a payload read is therefore one `EWSReadError`, one +disconnect, and no retry/error storm. + +## Protocol hardening retained inside the existing codec + +The patch does not replace the RFC codec, but it adds checks at its unsafe +boundaries: + +- canonical 16-bit and 64-bit payload lengths; +- rejection of the high bit in the RFC 63-bit length field; +- overflow-safe conversion to `SizeInt`; +- control frames must be final and no larger than 125 bytes; +- client endpoints reject masked server frames; +- server endpoints reject unmasked client frames; +- text messages and close reasons require valid UTF-8; +- close codes and outbound close payloads are validated; +- pong and close frames are delivered to `OnControl`; automatic pong and close + replies are written before application control callbacks run; +- an empty close is represented internally as 1005 and answered with an empty + close, never by putting reserved code 1005 on the wire; +- parser and close state are committed before callbacks; +- automatic protocol-close responses do not recursively invoke application + callbacks; +- no data frames may be sent once the close handshake has started. + +Incoming limits default to: + +| Property | Default | `0` means | +| --- | ---: | --- | +| `MaxFramePayloadSize` | 16 MiB | No policy limit; native `SizeInt` limits still apply. | +| `MaxMessagePayloadSize` | 64 MiB | No policy limit; native `SizeInt` limits still apply. | + +## Public compatibility decisions + +No existing public or protected declaration was removed or retyped, and +`IWSTransport` was not extended. The payload limits and typed exceptions are +additive. + +This is source compatibility, not PPU/ABI compatibility. Class layouts and +unit checksums changed. `fpwebsocket`, `fpwebsocketclient`, the companion units, +and the test must be rebuilt together; do not mix these sources with old +compiled `fpwebsocket*.ppu` files. + +There is one deliberate behavioral tightening: +`TWSMessagePump.AddClient(TWSClientConnection)` accepts only a +`TWebSocketClientConnection` carrying this unit's managed session. An arbitrary +base/custom connection has no reference, lease, or destruction-notification +contract. Letting the pump read it while an external owner can call `Free` +would preserve the old signature behavior only by retaining the old +use-after-free defect. + +The protected `CheckConnections` hook is still called before initial reader +scheduling. Once a persistent reader is running, returning `False` from a later +check is not a per-read pause/veto mechanism. + +## Known limitations + +- A concurrent client destructor cannot cancel a socket/TLS connect before the + new session is published. It can wait as long as `ConnectTimeout` and the + platform connect operation; the default zero value retains platform + semantics. +- The legacy TLS setup explicitly sets `VerifyPeerCert := False`; certificate + and hostname verification policy remains a separate security change. +- The default `Disconnect(True)` can block while writing its close frame. Only + reads currently have an explicit cross-thread cancellation primitive. +- The model uses one OS reader thread per live session, trading scalability for + simple blocking-I/O isolation and deterministic lifetime ownership. +- A non-threaded/custom pump which calls the inherited synchronous + `ReadConnections` still processes sessions sequentially and does not provide + stalled-sibling isolation. +- Destruction waits indefinitely for application callbacks that never return. + Two callbacks which synchronously wait while freeing each other's owners can + still deadlock. +- After a bounded `Terminate` returns while a synchronized callback is still + unwinding, an immediate `Execute` is a no-op until the old generation drains; + invoke `Execute` again afterwards. +- A raw `Connection` property value retained and used concurrently outside the + client API does not itself carry a session lease. +- Existing codec policies such as the configured/fixed outgoing mask key, + `woIndividualFrames`, and fragmented-message sequence flags were intentionally + left outside this bounded rewrite. + +## Companion server audit + +The unchanged `fpcustwsserver.pp`, `fpwebsocketserver.pp`, and `wsupgrader.pp` +use compatible shared `fpwebsocket` declarations. They do, however, need a +separate coherent lifetime pass: + +Server connections inherit the shared 16-MiB frame and 64-MiB message defaults, +but the companion server components do not currently expose properties for +configuring those limits. + +- `fpcustwsserver`: list locks span callbacks and blocking I/O; several worker + and pool paths retain raw connection pointers; removal and exceptional + cleanup are not consistently lifetime-safe. +- `fpwebsocketserver`: accept/handler shutdown can free a handler after an + unsuccessful wait; accept and connection-worker ownership needs one explicit + state machine. +- `wsupgrader`: `SetHost` assigns the property recursively; destruction does + not reliably unregister and quiesce; `StopUpgrader` ignores its wait result; + upgrade/handshake failure lacks ownership rollback. + +Those issues should not be fixed as unrelated one-line changes because the +server needs one ownership policy shared by accept, upgrade, worker, pool, and +connection-list code. + +## Building the regression test + +Run the compiler from the FPC source-tree root. The WebSocket source directory +must be an explicit unit path. A fresh output directory prevents an old +executable or PPU from being mistaken for the result of a failed compile. + +```bash +FPC_SOURCE_ROOT="$PWD" +BUILD_DIR="$(mktemp -d)" +FPC_BIN="${FPC_BIN:-fpc}" + +"$FPC_BIN" \ + -B -gl -vu \ + -Fu"$FPC_SOURCE_ROOT/packages/fcl-web/src/websocket" \ + -FE"$BUILD_DIR" \ + -FU"$BUILD_DIR" \ + "$FPC_SOURCE_ROOT/packages/fcl-web/tests/wsshutdowntest.pas" +``` + +`-B` forces the resolved WebSocket units and their dependent units to be +rebuilt. It does not choose which source copy wins; the explicit `-Fu` above +does that. The fresh `-FU` directory also prevents an old PPU from being reused +as output. + +For an auditable source-resolution log, replace `-vu` with `-va`, redirect the +compiler output to a file, and confirm that `fpwebsocket.pp` and +`fpwebsocketclient.pp` were loaded from +`packages/fcl-web/src/websocket`, not only as installed PPUs from the compiler +configuration. + +Run on Linux: + +```bash +"$BUILD_DIR/wsshutdowntest" +"$BUILD_DIR/wsshutdowntest" --pump-first +"$BUILD_DIR/wsshutdowntest" --interrupt-race +``` + +Use the corresponding `.exe` for Windows. The default run launches each of the +37 numbered scenarios in its own supervised process. TLS scenarios skip +themselves when no usable OpenSSL handler is available. + +The complete harness was run with FPC 3.3.1 on x86-64 Linux: all 37 scenarios, +`--pump-first`, and `--interrupt-race` passed, including plain-TCP and TLS +blocked-read termination. The same complete run remains to be repeated on +Windows for this rewrite. diff --git a/packages/fcl-web/src/websocket/fpwebsocket.pp b/packages/fcl-web/src/websocket/fpwebsocket.pp index cbf9db0c..1325842d 100644 --- a/packages/fcl-web/src/websocket/fpwebsocket.pp +++ b/packages/fcl-web/src/websocket/fpwebsocket.pp @@ -80,9 +80,39 @@ interface CLOSE_INTERNAL_SERVER_ERROR = 1011; CLOSE_TLS_HANDSHAKE = 1015; + { These limits apply to incoming data. Set the corresponding connection + property to 0 to disable the policy limit. The native TBytes/SizeInt + limit is always enforced. } + DefaultMaxFramePayloadSize = 16 * 1024 * 1024; + DefaultMaxMessagePayloadSize = 64 * 1024 * 1024; + type EWebSocket = Class(Exception); EWSHandShake = class(EWebSocket); + EWSReadInterrupted = class(EWebSocket); + + { Typed transport failures let a message pump distinguish an orderly EOF, + an interrupted read and a terminal socket error. } + EWSTransportError = class(EWebSocket) + private + FErrorCode: LongInt; + public + constructor Create(const aMessage: String; aErrorCode: LongInt); reintroduce; + property ErrorCode: LongInt read FErrorCode; + end; + EWSConnectionClosed = class(EWSTransportError); + EWSReadError = class(EWSTransportError); + EWSWriteError = class(EWSTransportError); + + { RFC protocol errors carry the close status which should be returned to + the peer before the connection is discarded. } + EWSProtocolError = class(EWebSocket) + private + FCloseCode: Word; + public + constructor Create(const aMessage: String; aCloseCode: Word); reintroduce; + property CloseCode: Word read FCloseCode; + end; TFrameType = (ftContinuation,ftText,ftBinary,ftClose,ftPing,ftPong,ftFutureOpcodes); @@ -94,7 +124,7 @@ EWSHandShake = class(EWebSocket); TIncomingResult = (irNone, // No data waiting irWaiting, // Data waiting irOK, // Data was waiting and handled - irClose // Data was waiting, handled, and we must disconnect (CloseState=csClosed) + irClose // The connection reached a terminal read/close condition ); { TFrameTypeHelper } @@ -194,8 +224,19 @@ TWSHeaders = class TWSSocketHelper = Class (TObject,IWSTransport) Private FSocket : TSocketStream; + FReadState : LongInt; + FReadsCancelled : LongInt; + FWriteLock : TRTLCriticalSection; + Procedure BeginRead; + Procedure EndRead; + function ReadSocket(var aBuffer; aCount : LongInt) : LongInt; + Procedure ReadSocketBuffer(var aBuffer; aCount : LongInt); + function WriteSocket(const aBuffer; aCount : LongInt) : LongInt; Public Constructor Create (aSocket : TSocketStream); + Destructor Destroy; override; + Procedure InterruptRead; + Procedure CancelReads; Function CanRead(aTimeOut: Integer) : Boolean; function PeerIP: string; virtual; function PeerPort: word; virtual; @@ -216,6 +257,8 @@ TWSTransport = class(TObject, IWSTransport) Constructor Create(aStream : TSocketStream); Destructor Destroy; override; Procedure CloseSocket; + Procedure InterruptRead; + Procedure CancelReads; Property Helper : TWSSocketHelper Read FHelper Implements IWSTransport; Property Socket : TSocketStream Read GetSocket; end; @@ -230,7 +273,9 @@ TWSFramePayload = record MaskKey: dword; Masked: Boolean; Procedure ReadData(var Content : TBytes; aTransport : IWSTransport); - Procedure Read(buffer: TBytes; aTransport : IWSTransport); + Procedure Read(buffer: TBytes; aTransport : IWSTransport); overload; + Procedure Read(buffer: TBytes; aTransport : IWSTransport; + aMaxPayloadSize: QWord); overload; class procedure DoMask(var aData: TBytes; Key: DWORD); static; class procedure CopyMasked(SrcData: TBytes; var DestData: TBytes; Key: DWORD; aOffset: Integer); static; class function CopyMasked(SrcData: TBytes; Key: DWORD) : TBytes; static; @@ -245,11 +290,15 @@ TWSFramePayload = record FPayload : TWSFramePayload; FReason: WORD; protected - function Read(aTransport: IWSTransport): boolean; + function Read(aTransport: IWSTransport): boolean; overload; + function Read(aTransport: IWSTransport; + aMaxPayloadSize: QWord): boolean; overload; function GetAsBytes : TBytes; virtual; Public // Read a message from transport. Returns Nil if the connection was closed when reading. - class function CreateFromStream(aTransport : IWSTransport): TWSFrame; + class function CreateFromStream(aTransport : IWSTransport): TWSFrame; overload; + class function CreateFromStream(aTransport : IWSTransport; + aMaxPayloadSize: QWord): TWSFrame; overload; public constructor Create(aType: TFrameType; aIsFinal: Boolean; APayload: TBytes; aMask : Integer = 0); overload; virtual; constructor Create(Const aMessage : UTF8String; aMask : Integer = 0); overload; virtual; @@ -321,8 +370,19 @@ TWSConnection = class FOnControl: TWSControlEvent; FCloseState : TCloseState; FOptions: TWSOptions; + FMaxFramePayloadSize: QWord; + FMaxMessagePayloadSize: QWord; + FWriteLock: TRTLCriticalSection; Function GetPeerIP : String; Function GetPeerPort : word; + function GetCloseState: TCloseState; + function MessageSizeAllowed(aCurrentSize, aAdditionalSize: SizeInt): Boolean; + procedure PrepareProtocolClose(out aSendClose: Boolean); + function WriteProtocolClose(aReason: Word): Boolean; + function WriteFrame(aFrame: TWSFrame; + out aErrorMessage: String): Boolean; + function WriteControlFrame(aFrameType: TFrameType; const aData: TBytes; + out aErrorMessage: String): Boolean; protected procedure AllocateConnectionID; virtual; Procedure SetCloseState(aValue : TCloseState); virtual; @@ -368,7 +428,7 @@ TWSConnection = class // Disconnect when status is set to csClosed; Property AutoDisconnect : Boolean Read FAutoDisconnect Write FAutoDisconnect; // Close frame handling - Property CloseState : TCloseState Read FCloseState; + Property CloseState : TCloseState Read GetCloseState; // Connection ID, allocated during create Property ConnectionID : String Read FConnectionID; // If set to true, the owner data is freed when the connection is freed. @@ -381,6 +441,10 @@ TWSConnection = class Property Options : TWSOptions Read FOptions; // Mask to use when sending frames. Set to nonzero value to send masked frames. Property OutgoingFrameMask : Integer Read FOutgoingFrameMask Write FOutgoingFrameMask; + // Maximum accepted payload of one incoming frame. 0 disables the policy limit. + property MaxFramePayloadSize: QWord read FMaxFramePayloadSize write FMaxFramePayloadSize; + // Maximum accepted assembled incoming data message. 0 disables the policy limit. + property MaxMessagePayloadSize: QWord read FMaxMessagePayloadSize write FMaxMessagePayloadSize; // Peer IP address property PeerIP: string read GetPeerIP; // Peer IP port @@ -419,7 +483,7 @@ TWSConnection = class function GetHandshakeCompleted: Boolean; override; // Owned by connection Property ClientTransport : TWSClientTransport Read FTransport; - // + // Non-owning compatibility alias; the assigning code retains ownership. Property HandShakeResponse : TWSHandShakeResponse Read FHandshakeResponse Write FHandshakeResponse; End; @@ -491,6 +555,23 @@ TWSServerTransport = class(TWSTransport) SErrInvalidSizeFlag = 'Invalid size flag: %d'; SErrInvalidFrameType = 'Invalid frame type flag: %d'; SErrWriteReturnedError = 'Write operation returned error: (%d) %s'; + SErrWriteClosed = 'Write operation returned no data; the connection is closed'; + SErrReadReturnedError = 'Read operation returned error: (%d) %s'; + SErrReadClosed = 'The peer closed the WebSocket connection'; + SErrInvalidReadCount = 'Transport returned invalid read count %d for a %d-byte request'; + SErrInvalidWriteCount = 'Transport returned invalid write count %d for a %d-byte request'; + SErrNegativeBufferSize = 'Negative WebSocket buffer size: %d'; + SErrHandshakeLineTooLong = 'WebSocket handshake line exceeds the transport limit'; + SErrHandshakeHeadersTooLarge = 'WebSocket handshake headers exceed the transport limit'; + SErrNonCanonicalLength = 'Non-canonical WebSocket payload length encoding'; + SErrInvalidPayloadLength = 'Invalid WebSocket payload length'; + SErrFramePayloadTooLarge = 'WebSocket frame payload is too large (%s bytes)'; + SErrControlPayloadTooLarge = 'WebSocket control frame payload exceeds 125 bytes'; + SErrPayloadLengthMismatch = 'WebSocket payload length does not match its data buffer'; + SErrInvalidFrameMask = 'Invalid WebSocket frame masking for this endpoint'; + SErrConnectionClosing = 'Cannot send WebSocket data after the close handshake has started'; + SErrReadInterrupted = 'WebSocket read interrupted'; + SErrConcurrentRead = 'Concurrent reads on one WebSocket transport are not supported'; function DecodeBytesBase64(const s: string; Strict: boolean = false) : TBytes; function EncodeBytesBase64(const aBytes : TBytes) : String; @@ -504,6 +585,135 @@ implementation uses strutils, sha1, base64; {$ENDIF FPC_DOTTEDUNITS} +Const + WSReadIdle = 0; + WSReadActive = 1; + WSReadInterrupting = 2; + WSMaxTransportLineBytes = 64 * 1024; + WSMaxServerHandshakeLineBytes = 8192; + WSMaxServerHandshakeLines = 256; + WSMaxServerHandshakeBytes = 64 * 1024; + +{ EWSTransportError } + +constructor EWSTransportError.Create(const aMessage: String; + aErrorCode: LongInt); +begin + FErrorCode:=aErrorCode; + inherited Create(aMessage); +end; + +{ EWSProtocolError } + +constructor EWSProtocolError.Create(const aMessage: String; aCloseCode: Word); +begin + FCloseCode:=aCloseCode; + inherited Create(aMessage); +end; + +function WSLastSocketError: LongInt; inline; +begin + Result:={$IFDEF FPC_DOTTEDUNITS}System.Net.{$ENDIF}sockets.SocketError; +end; + +function WSValidCloseCode(aCode: Word): Boolean; inline; +begin + Result:=((aCode>=CLOSE_NORMAL_CLOSURE) and + (aCode<=1014) and + (aCode<>CLOSE_RESERVER) and + (aCode<>CLOSE_NO_STATUS_RCVD) and + (aCode<>CLOSE_ABNORMAL_CLOSURE)) or + ((aCode>=3000) and (aCode<=4999)); +end; + +procedure WSRaiseReadError; +var + ErrorCode: LongInt; +begin + { SocketError must be sampled before formatting or allocating anything. } + ErrorCode:=WSLastSocketError; + raise EWSReadError.Create( + Format(SErrReadReturnedError,[ErrorCode,SysErrorMessage(ErrorCode)]), + ErrorCode); +end; + +procedure WSReadExact(aTransport: IWSTransport; var aBytes: TBytes; + aCount: Integer); +var + Chunk: TBytes; + BytesRead: Integer; + ReadPosition: Integer; + Remaining: Integer; +begin + if aCount<0 then + raise EWebSocket.CreateFmt(SErrNegativeBufferSize,[aCount]); + SetLength(aBytes,aCount); + if aCount=0 then + Exit; + if not Assigned(aTransport) then + raise EWSConnectionClosed.Create(SErrReadClosed,0); + + ReadPosition:=0; + while ReadPositionRemaining) or (Length(Chunk)0 then + Raise EWSReadInterrupted.Create(SErrReadInterrupted); + PreviousState:=InterlockedCompareExchange(FReadState,WSReadActive, + WSReadIdle); + if PreviousState<>WSReadIdle then + begin + if (PreviousState=WSReadInterrupting) or + (InterlockedCompareExchange(FReadsCancelled,0,0)<>0) then + Raise EWSReadInterrupted.Create(SErrReadInterrupted); + Raise EWebSocket.Create(SErrConcurrentRead); + end; + { Close the race in which terminal cancellation is published immediately + after the first check but before this read claims FReadState. } + if InterlockedCompareExchange(FReadsCancelled,0,0)<>0 then + begin + InterlockedCompareExchange(FReadState,WSReadIdle,WSReadActive); + Raise EWSReadInterrupted.Create(SErrReadInterrupted); + end; +end; + +procedure TWSSocketHelper.EndRead; +Var + PreviousState : LongInt; +begin + { The state transition made by InterruptRead is itself the interruption + publication. This single atomic exchange closes the former window between + claiming a read and publishing a separate interruption flag. } + PreviousState:=InterlockedExchange(FReadState,WSReadIdle); + if (PreviousState=WSReadInterrupting) or + (InterlockedCompareExchange(FReadsCancelled,0,0)<>0) then + Raise EWSReadInterrupted.Create(SErrReadInterrupted); +end; + +function TWSSocketHelper.ReadSocket(var aBuffer; aCount: LongInt): LongInt; +var + ErrorCode: LongInt; + ReadResult: LongInt; +begin + ErrorCode:=0; + BeginRead; + try + ReadResult:=FSocket.Read(aBuffer,aCount); + if ReadResult<0 then + { Capture this before EndRead performs any synchronization operation. } + ErrorCode:=WSLastSocketError; + finally + EndRead; + end; + if ReadResult<0 then + raise EWSReadError.Create( + Format(SErrReadReturnedError, + [ErrorCode,SysErrorMessage(ErrorCode)]),ErrorCode); + Result:=ReadResult; +end; + +procedure TWSSocketHelper.ReadSocketBuffer(var aBuffer; aCount: LongInt); +var + Buffer: TBytes; + BytesRead: LongInt; + ReadPosition: LongInt; +begin + if aCount<0 then + raise EWebSocket.CreateFmt(SErrNegativeBufferSize,[aCount]); + if aCount=0 then + Exit; + SetLength(Buffer,aCount); + ReadPosition:=0; + while ReadPositionaCount-ReadPosition then + raise EWSReadError.Create( + Format(SErrInvalidReadCount,[BytesRead,aCount-ReadPosition]),0); + Inc(ReadPosition,BytesRead); + end; + Move(Buffer[0],aBuffer,aCount); +end; + +function TWSSocketHelper.WriteSocket(const aBuffer; aCount: LongInt): LongInt; +var + ErrorCode: LongInt; +begin + Result:=FSocket.Write(aBuffer,aCount); + if Result<0 then + begin + { Capture this before releasing the surrounding write lock. } + ErrorCode:=WSLastSocketError; + raise EWSWriteError.Create( + Format(SErrWriteReturnedError, + [ErrorCode,SysErrorMessage(ErrorCode)]),ErrorCode); + end; +end; + +procedure TWSSocketHelper.InterruptRead; +Var + PreviousState : LongInt; +begin + { Claim only a transport whose reader is currently inside a socket read. + A pump may serve several connections, so interrupting every registered + socket would unnecessarily break healthy siblings. } + PreviousState:=InterlockedCompareExchange(FReadState,WSReadInterrupting, + WSReadActive); + if (PreviousState<>WSReadActive) and + (PreviousState<>WSReadInterrupting) then + Exit; + + { Keep WSReadInterrupting set until EndRead observes it and raises. Repeated + termination passes may retry the platform wake, but the connection can no + longer return to the pump as healthy after its socket has been shut down. } + WakeSocketRead(FSocket.Handle); +end; + +procedure TWSSocketHelper.CancelReads; +begin + { Unlike InterruptRead, this is a terminal per-session cancellation. It is + used only after the owning session has been deregistered, and deliberately + prevents a reader from entering a later exact-read chunk. } + InterlockedExchange(FReadsCancelled,1); + InterlockedExchange(FReadState,WSReadInterrupting); + { Shutting down while idle makes the cancellation durable at the socket + layer too. The descriptor remains owned by TWSTransport. } + WakeSocketRead(FSocket.Handle); +end; + function TWSSocketHelper.CanRead(aTimeOut: Integer): Boolean; begin Result:=FSocket.CanRead(aTimeout); @@ -650,6 +1017,7 @@ function TWSSocketHelper.ReadLn: String; Var C : Byte; + BytesRead, aSize : integer; begin @@ -658,55 +1026,94 @@ function TWSSocketHelper.ReadLn: String; SetLength(Result,255); aSize:=0; C:=0; - While (FSocket.Read(C,1)=1) and (C<>10) do - begin + repeat + BytesRead:=ReadSocket(C,1); + if BytesRead=0 then + raise EWSConnectionClosed.Create(SErrReadClosed,0); + if BytesRead<>1 then + raise EWSReadError.Create( + Format(SErrInvalidReadCount,[BytesRead,1]),0); + if C=10 then + Break; + if aSize>=WSMaxTransportLineBytes then + raise EWSHandShake.Create(SErrHandshakeLineTooLong); Inc(aSize); if aSize>Length(Result) then SetLength(Result,Length(Result)+255); Result[aSize]:=AnsiChar(C); - end; + until False; if (aSize>0) and (Result[aSize]=#13) then Dec(aSize); SetLength(Result,aSize); end; function TWSSocketHelper.ReadBytes(var aBytes: TBytes; aCount: Integer): Integer; -var - buf: TBytes; - aPos, toRead: QWord; begin - if aCount=0 then exit(0); - aPos := 0; - SetLength(aBytes, aCount); - repeat - SetLength(buf{%H-}, aCount); - Result := FSocket.Read(buf[0], aCount - aPos); - if Result <= 0 then - break; - SetLength(buf, Result); - Move(buf[0], aBytes[aPos], Result); - Inc(aPos, Result); - ToRead := aCount - aPos; - Result := aCount; - until toRead <= 0; + if aCount<0 then + raise EWebSocket.CreateFmt(SErrNegativeBufferSize,[aCount]); + SetLength(aBytes,aCount); + if aCount=0 then + Exit(0); + Result:=ReadSocket(aBytes[0],aCount); + if Result>aCount then + raise EWSReadError.Create( + Format(SErrInvalidReadCount,[Result,aCount]),0); + SetLength(aBytes,Result); end; procedure TWSSocketHelper.ReadBuffer(aBytes: TBytes); begin if Length(ABytes)=0 then exit; - FSocket.ReadBuffer(aBytes[0],Length(ABytes)); + ReadSocketBuffer(aBytes[0],Length(ABytes)); end; function TWSSocketHelper.WriteBytes(aBytes: TBytes; aCount: Integer): Integer; begin - if aCount=0 then exit(0); - Result:=FSocket.Write(aBytes[0],aCount); + if (aCount<0) or (aCount>Length(aBytes)) then + raise EWSWriteError.Create( + Format(SErrInvalidWriteCount,[aCount,Length(aBytes)]),0); + if aCount=0 then + Exit(0); + EnterCriticalSection(FWriteLock); + try + Result:=WriteSocket(aBytes[0],aCount); + if Result>aCount then + raise EWSWriteError.Create( + Format(SErrInvalidWriteCount,[Result,aCount]),0); + finally + LeaveCriticalSection(FWriteLock); + end; end; procedure TWSSocketHelper.WriteBuffer(aBytes: TBytes); +var + BytesWritten: LongInt; + WriteCount: LongInt; + WritePosition: SizeInt; begin - if Length(aBytes)=0 then exit; - FSocket.WriteBuffer(aBytes[0],Length(aBytes)); + if Length(aBytes)=0 then + Exit; + EnterCriticalSection(FWriteLock); + try + WritePosition:=0; + while WritePositionHigh(LongInt) then + WriteCount:=High(LongInt) + else + WriteCount:=LongInt(Length(aBytes)-WritePosition); + BytesWritten:=WriteSocket(aBytes[WritePosition],WriteCount); + if BytesWritten=0 then + raise EWSWriteError.Create(SErrWriteClosed,0); + if BytesWritten>WriteCount then + raise EWSWriteError.Create( + Format(SErrInvalidWriteCount, + [BytesWritten,WriteCount]),0); + Inc(WritePosition,BytesWritten); + end; + finally + LeaveCriticalSection(FWriteLock); + end; end; { TWSMessage } @@ -891,73 +1298,93 @@ procedure TWSFramePayload.ReadData(var Content: TBytes; aTransport: IWSTransport Var Buf : TBytes; - aPos,toRead : QWord; - aCount, FailCnt : Longint; + ReadPosition, + ToRead : QWord; + ReadCount : Longint; begin Buf:=[]; ToRead:=DataLength; - aPos:=0; - FailCnt:=0; - Repeat - aCount:=ToRead; - if aCount>MaxBufSize then - aCount:=MaxBufSize; - SetLength(Buf,aCount); - aCount := aTransport.ReadBytes(Buf,aCount); - if aCount>0 then - begin - Move(Buf[0],Content[aPos],aCount); - Inc(aPos,aCount); - ToRead:=DataLength-aPos; - FailCnt:=0; - end + ReadPosition:=0; + while ToRead>0 do + begin + if ToRead>MaxBufSize then + ReadCount:=MaxBufSize else - begin - sleep(1); - inc(FailCnt); - if FailCnt>100 then - raise Exception.Create('20230316102741 TWSFramePayload.ReadData'); - end; - Until (ToRead<=0); + ReadCount:=LongInt(ToRead); + WSReadExact(aTransport,Buf,ReadCount); + Move(Buf[0],Content[SizeInt(ReadPosition)],ReadCount); + Inc(ReadPosition,QWord(ReadCount)); + Dec(ToRead,QWord(ReadCount)); + end; end; procedure TWSFramePayload.Read(buffer: TBytes; aTransport: IWSTransport); +begin + Read(buffer,aTransport,0); +end; + +procedure TWSFramePayload.Read(buffer: TBytes; aTransport: IWSTransport; + aMaxPayloadSize: QWord); Var LenFlag : Byte; - paylen16 : Word; - content: TBytes; + PayLen16 : Word; + EncodedLength, + Content: TBytes; begin + if Length(Buffer)<2 then + raise EWSProtocolError.Create(SErrInvalidPayloadLength, + CLOSE_PROTOCOL_ERROR); content:=[]; + Data:=[]; + MaskKey:=0; Masked := ((buffer[1] and FlagMasked) <> 0); LenFlag := buffer[1] and FlagLengthMask; Case LenFlag of FlagTwoBytes: begin - aTransport.ReadBytes(Buffer,2); - Paylen16:=Buffer.ToWord(0); + WSReadExact(aTransport,EncodedLength,2); + Paylen16:=EncodedLength.ToWord(0); DataLength := ntohs(PayLen16); + if DataLength0 then + raise EWSProtocolError.Create(SErrInvalidPayloadLength, + CLOSE_PROTOCOL_ERROR); + DataLength:=EncodedLength.ToQWord(0); + DataLength:=ntohx(DataLength); + if DataLength<(QWord(1) shl 16) then + raise EWSProtocolError.Create(SErrNonCanonicalLength, + CLOSE_PROTOCOL_ERROR); end else DataLength:=lenFlag; end; + if (aMaxPayloadSize<>0) and (DataLength>aMaxPayloadSize) then + raise EWSProtocolError.Create( + Format(SErrFramePayloadTooLarge,[IntToStr(Int64(DataLength))]), + CLOSE_MESSAGE_TOO_BIG); + if DataLength>QWord(High(SizeInt)) then + raise EWSProtocolError.Create( + Format(SErrFramePayloadTooLarge,[IntToStr(Int64(DataLength))]), + CLOSE_MESSAGE_TOO_BIG); + if Masked then - begin - // In some times, not 4 bytes are returned - aTransport.ReadBytes(Buffer,4); - MaskKey:=buffer.ToDword(0); - end; - SetLength(content, DataLength); + begin + WSReadExact(aTransport,EncodedLength,4); + MaskKey:=EncodedLength.ToDword(0); + end; + SetLength(content,SizeInt(DataLength)); if (DataLength>0) then begin ReadData(Content,aTransport); @@ -977,7 +1404,7 @@ constructor TWSFrame.Create(aType: TFrameType; aIsFinal: Boolean; APayload: TByt FPayload.Data := APayload; if Assigned(aPayload) then - FPayload.DataLength := Cardinal(Length(aPayload)); + FPayload.DataLength := QWord(Length(aPayload)); end; constructor TWSFrame.Create(aType: TFrameType; aIsFinal : Boolean = True; aMask: Integer=0); @@ -1002,11 +1429,16 @@ constructor TWSFrame.Create(const aMessage: UTF8String; aMask: Integer=0); end; class function TWSFrame.CreateFromStream(aTransport : IWSTransport): TWSFrame; +begin + Result:=CreateFromStream(aTransport,0); +end; +class function TWSFrame.CreateFromStream(aTransport : IWSTransport; + aMaxPayloadSize: QWord): TWSFrame; begin Result:=TWSFrame.Create; try - if not Result.Read(aTransport) then + if not Result.Read(aTransport,aMaxPayloadSize) then FreeAndNil(Result); except FreeAndNil(Result); @@ -1016,36 +1448,64 @@ class function TWSFrame.CreateFromStream(aTransport : IWSTransport): TWSFrame; function TWSFrame.Read(aTransport: IWSTransport): boolean; +begin + Result:=Read(aTransport,0); +end; +function TWSFrame.Read(aTransport: IWSTransport; + aMaxPayloadSize: QWord): boolean; Var Buffer : Tbytes; - B1 : Byte; + B1, + LenFlag : Byte; begin Result:=False; Buffer:=Default(TBytes); - SetLength(Buffer,2); - if aTransport.ReadBytes(Buffer,2)=0 then - Exit; - if Length(Buffer)<2 then - Raise EWebSocket.Create('Could not read frame header'); + try + WSReadExact(aTransport,Buffer,2); + except + { EOF before a new frame is the normal terminal indication expected by + CheckIncoming. EOF after a partial fixed field is also terminal and the + incomplete frame is never exposed. } + on E: EWSConnectionClosed do + Exit; + end; B1:=buffer[0]; FFinalFrame:=(B1 and FlagFinalFrame) = FlagFinalFrame; FRSV:=(B1 and %01110000) shr 4; FFrameType.AsFlag:=(B1 and $F); - FPayload.Read(Buffer,aTransport); - FReason:=CLOSE_NORMAL_CLOSURE; + if (FRSV<>0) or (FFrameType=ftFutureOpcodes) then + raise EWSProtocolError.Create( + Format(SErrInvalidFrameType,[B1 and $F]),CLOSE_PROTOCOL_ERROR); + LenFlag:=Buffer[1] and FlagLengthMask; + if FFrameType in [ftClose,ftPing,ftPong] then + begin + if not FFinalFrame then + raise EWSProtocolError.Create( + Format(SErrInvalidFrameType,[B1 and $F]),CLOSE_PROTOCOL_ERROR); + if LenFlag>125 then + raise EWSProtocolError.Create(SErrControlPayloadTooLarge, + CLOSE_PROTOCOL_ERROR); + end; + FPayload.Read(Buffer,aTransport,aMaxPayloadSize); + if FFrameType=ftClose then + FReason:=CLOSE_NO_STATUS_RCVD + else + FReason:=CLOSE_NORMAL_CLOSURE; if FFrameType=ftClose then if FPayload.DataLength = 1 then - FReason:=CLOSE_PROTOCOL_ERROR + raise EWSProtocolError.Create(SErrInvalidPayloadLength, + CLOSE_PROTOCOL_ERROR) else if FPayload.DataLength>1 then begin - FReason:=SwapEndian(FPayload.Data.ToWord(0)); + FReason:=NToHs(FPayload.Data.ToWord(0)); FPayload.DataLength := FPayload.DataLength - 2; if FPayload.DataLength > 0 then - move(FPayload.Data[2], FPayload.Data[0], FPayload.DataLength); - SetLength(FPayload.Data, FPayload.DataLength); + move(FPayload.Data[2],FPayload.Data[0], + SizeInt(FPayload.DataLength)); + SetLength(FPayload.Data,SizeInt(FPayload.DataLength)); end; Result:=True; end; @@ -1062,6 +1522,17 @@ function TWSFrame.GetAsBytes: TBytes; begin Result:=Nil; + if FPayload.DataLength<>QWord(Length(FPayload.Data)) then + raise EWSProtocolError.Create(SErrPayloadLengthMismatch, + CLOSE_PROTOCOL_ERROR); + if (FrameType in [ftClose,ftPing,ftPong]) and + (FPayload.DataLength>125) then + raise EWSProtocolError.Create(SErrControlPayloadTooLarge, + CLOSE_PROTOCOL_ERROR); + if FPayload.DataLength>QWord(High(SizeInt)-14) then + raise EWSProtocolError.Create( + Format(SErrFramePayloadTooLarge, + [IntToStr(Int64(FPayload.DataLength))]),CLOSE_MESSAGE_TOO_BIG); firstByte := FrameType.AsFlag; if FinalFrame then firstByte := firstByte or FlagFinalFrame; @@ -1094,7 +1565,7 @@ function TWSFrame.GetAsBytes: TBytes; lenByte:=Lenbyte or FlagMasked; aoffSet:=aOffSet+4; end; - SetLength(buffer,aOffset+Int64(FPayload.DataLength)); + SetLength(buffer,aOffset+SizeInt(FPayload.DataLength)); buffer[0] := firstByte; buffer[1] := LenByte; for I := 0 to Length(LengthBytes)-1 do @@ -1106,7 +1577,8 @@ function TWSFrame.GetAsBytes: TBytes; end else if Payload.DataLength > 0 then - move(Payload.Data[0], buffer[aOffset], Payload.DataLength); + move(Payload.Data[0],buffer[aOffset], + SizeInt(Payload.DataLength)); Result := Buffer; end; @@ -1120,9 +1592,9 @@ class procedure TWSFramePayload.DoMask(var aData: TBytes; Key: DWORD); class procedure TWSFramePayload.CopyMasked(SrcData: TBytes; var DestData: TBytes; Key: DWORD; aOffset: Integer); var - currentMaskIndex: Longint; + currentMaskIndex: SizeInt; byteKeys: TBytes; - I: Longint; + I: SizeInt; begin CurrentMaskIndex := 0; @@ -1216,18 +1688,30 @@ procedure TWSConnection.SetHandShakeRequest(aRequest: TWSHandShakeRequest); constructor TWSConnection.Create(aOwner : TComponent; aOptions: TWSOptions); begin + InitCriticalSection(FWriteLock); FOwner:=aOwner; Foptions:=aOptions; - FWebSocketVersion:=WebSocketVersion; + FWebSocketVersion:=DefaultWebSocketVersion; + FInitialOpcode:=ftContinuation; + FCloseState:=csNone; + FMaxFramePayloadSize:=DefaultMaxFramePayloadSize; + FMaxMessagePayloadSize:=DefaultMaxMessagePayloadSize; AllocateConnectionID; end; destructor TWSConnection.Destroy; begin - FreeAndNil(FHandshakeRequest); - If FreeUserData then - FreeAndNil(FUserData); - inherited; + try + FreeAndNil(FHandshakeRequest); + finally + try + If FreeUserData then + FreeAndNil(FUserData); + finally + DoneCriticalSection(FWriteLock); + inherited Destroy; + end; + end; end; class function TWSConnection.GetCloseData(aBytes: TBytes; out aReason: String): Word; @@ -1277,12 +1761,86 @@ procedure TWSConnection.AllocateConnectionID; end; procedure TWSConnection.SetCloseState(aValue: TCloseState); +var + MustDisconnect: Boolean; begin - FCloseState:=aValue; - if (FCloseState=csClosed) and AutoDisconnect then + EnterCriticalSection(FWriteLock); + try + FCloseState:=aValue; + MustDisconnect:=(FCloseState=csClosed) and AutoDisconnect; + finally + LeaveCriticalSection(FWriteLock); + end; + if MustDisconnect then Disconnect; end; +function TWSConnection.GetCloseState: TCloseState; +begin + EnterCriticalSection(FWriteLock); + try + Result:=FCloseState; + finally + LeaveCriticalSection(FWriteLock); + end; +end; + +function TWSConnection.MessageSizeAllowed(aCurrentSize, + aAdditionalSize: SizeInt): Boolean; +var + TotalSize: QWord; +begin + if (aCurrentSize<0) or (aAdditionalSize<0) then + Exit(False); + TotalSize:=QWord(aCurrentSize)+QWord(aAdditionalSize); + Result:=(TotalSize<=QWord(High(SizeInt))) and + ((FMaxMessagePayloadSize=0) or + (TotalSize<=FMaxMessagePayloadSize)); +end; + +procedure TWSConnection.PrepareProtocolClose(out aSendClose: Boolean); +begin + EnterCriticalSection(FWriteLock); + try + case FCloseState of + csNone: + begin + FCloseState:=csReceived; + aSendClose:=True; + end; + csReceived: + aSendClose:=False; + csSent: + begin + FCloseState:=csClosed; + aSendClose:=False; + end; + csClosed: + aSendClose:=False; + end; + finally + LeaveCriticalSection(FWriteLock); + end; +end; + +function TWSConnection.WriteProtocolClose(aReason: Word): Boolean; +var + CloseData: TBytes; + ErrorMessage: String; + SendClose: Boolean; +begin + PrepareProtocolClose(SendClose); + if not SendClose then + Exit(False); + SetLength(CloseData,2); + CloseData[0]:=(aReason and $FF00) shr 8; + CloseData[1]:=aReason and $FF; + { Protocol-generated close frames never dispatch a nested application + callback. The pump reports any terminal read/write failure once, after + the connection has been retired. } + Result:=WriteControlFrame(ftClose,CloseData,ErrorMessage); +end; + function TWSConnection.ReadMessage: Boolean; begin Result:=DoReadMessage; @@ -1292,29 +1850,35 @@ procedure TWSConnection.DispatchEvent(aInitialType: TFrameType; aFrame: TWSFrame Var msg: TWSMessage; + ControlHandler: TWSControlEvent; + MessageHandler: TWSMessageEvent; begin Case aInitialType of ftPing, ftPong, ftClose : - If Assigned(FOnControl) then - FOnControl(Self,aInitialType,aMessageContent); + begin + ControlHandler:=FOnControl; + If Assigned(ControlHandler) then + ControlHandler(Self,aInitialType,aMessageContent); + end; ftBinary, ftText : begin - if Assigned(FOnMessageReceived) then + MessageHandler:=FOnMessageReceived; + if Assigned(MessageHandler) then begin Msg:=Default(TWSMessage); Msg.IsText:=(aInitialType=ftText); - if aFrame.FrameType=ftBinary then + if aFrame.FrameType in [ftBinary,ftText] then Msg.Sequences:=[fsFirst] else Msg.Sequences:=[fsContinuation]; if aFrame.FinalFrame then Msg.Sequences:=Msg.Sequences+[fsLast]; Msg.PayLoad:=aMessageContent; - FOnMessageReceived(Self, Msg); + MessageHandler(Self, Msg); end; end; ftContinuation: ; // Cannot happen normally @@ -1323,152 +1887,191 @@ procedure TWSConnection.DispatchEvent(aInitialType: TFrameType; aFrame: TWSFrame function TWSConnection.HandleIncoming(aFrame: TWSFrame) : Boolean; - Procedure UpdateCloseState; + procedure ProtocolError(aCode: Word); + begin + { Commit the terminal state before sending. } + Result:=False; + WriteProtocolClose(aCode); + end; - begin - if (FCloseState=csNone) then - FCloseState:=csReceived - else if (FCloseState=csSent) then - FCloseState:=csClosed; - end; + function ReceiveClose: TCloseState; + begin + EnterCriticalSection(FWriteLock); + try + Result:=FCloseState; + case FCloseState of + csNone: + FCloseState:=csReceived; + csSent: + FCloseState:=csClosed; + csReceived, + csClosed: + ; + end; + finally + LeaveCriticalSection(FWriteLock); + end; + end; - procedure ProtocolError(aCode: Word); - begin - Close('', aCode); - UpdateCloseState; - Result:=false; - end; +var + CloseData, + ErrorData: TBytes; + WriteErrorMessage: String; + InitialType: TFrameType; + MessageContent: TBytes; + PreviousCloseState: TCloseState; begin - Result:=true; + Result:=True; // check Reserved bits if aFrame.Reserved<>0 then - begin + begin ProtocolError(CLOSE_PROTOCOL_ERROR); Exit; - end; + end; // check Reserved opcode - if aFrame.FrameType = ftFutureOpcodes then - begin + if aFrame.FrameType=ftFutureOpcodes then + begin ProtocolError(CLOSE_PROTOCOL_ERROR); Exit; - end; - { If control frame it must be complete } - if ((aFrame.FrameType=ftPing) or - (aFrame.FrameType=ftPong) or - (aFrame.FrameType=ftClose)) - and (not aFrame.FinalFrame) then - begin + end; + { A control frame must be final and have at most 125 payload bytes. } + if (aFrame.FrameType in [ftPing,ftPong,ftClose]) and + ((not aFrame.FinalFrame) or (aFrame.Payload.DataLength>125)) then + begin ProtocolError(CLOSE_PROTOCOL_ERROR); Exit; - end; - // - - // here we handle payload. -// if aFrame.FrameType in [ftBinary,ftText] then -// begin -// FInitialOpcode:=aFrame.FrameType; -// FMessageContent:=aFrame.Payload.Data; -// end; + end; - // Special handling Case aFrame.FrameType of ftContinuation: begin - if FInitialOpcode=ftContinuation then + if FInitialOpcode=ftContinuation then begin - ProtocolError(CLOSE_PROTOCOL_ERROR); - Exit; + ProtocolError(CLOSE_PROTOCOL_ERROR); + Exit; + end; + if not MessageSizeAllowed(Length(FMessageContent), + Length(aFrame.Payload.Data)) then + begin + ProtocolError(CLOSE_MESSAGE_TOO_BIG); + Exit; end; - FMessageContent.Append(aFrame.Payload.Data); - if aFrame.FinalFrame then + FMessageContent.Append(aFrame.Payload.Data); + if aFrame.FinalFrame then begin - if FInitialOpcode = ftText then - if IsValidUTF8(FMessageContent) then - DispatchEvent(FInitialOpcode,aFrame,FMessageContent) - else - ProtocolError(CLOSE_INVALID_FRAME_PAYLOAD_DATA) - else - DispatchEvent(FInitialOpcode,aFrame,FMessageContent); - // reset initial opcode - FInitialOpcode:=ftContinuation; + InitialType:=FInitialOpcode; + MessageContent:=FMessageContent; + { All parser state is reset before application code is entered. } + FInitialOpcode:=ftContinuation; + FMessageContent:=[]; + if (InitialType=ftText) and + (not IsValidUTF8(MessageContent)) then + begin + ProtocolError(CLOSE_INVALID_FRAME_PAYLOAD_DATA); + Exit; + end; + DispatchEvent(InitialType,aFrame,MessageContent); end; end; ftPing: begin - if aFrame.Payload.DataLength > 125 then - ProtocolError(CLOSE_PROTOCOL_ERROR) + if not (woPongExplicit in Options) then + if not WriteControlFrame(ftPong,aFrame.Payload.Data, + WriteErrorMessage) then + begin + Result:=False; + ErrorData:=TEncoding.UTF8.GetBytes( + UnicodeString(WriteErrorMessage)); + DispatchEvent(ftClose,nil,ErrorData); + Exit; + end; + { The callback is deliberately last: it may disconnect or free Self. } + DispatchEvent(ftPing,aFrame,aFrame.Payload.Data); + end; + + ftPong: + { The callback is deliberately last: it may disconnect or free Self. } + DispatchEvent(ftPong,aFrame,aFrame.Payload.Data); + + ftClose: + begin + if (aFrame.Payload.DataLength>123) or + ((aFrame.Reason<>CLOSE_NO_STATUS_RCVD) and + (not WSValidCloseCode(aFrame.Reason))) then + begin + ProtocolError(CLOSE_PROTOCOL_ERROR); + Exit; + end; + if not IsValidUTF8(aFrame.Payload.Data) then + begin + ProtocolError(CLOSE_INVALID_FRAME_PAYLOAD_DATA); + Exit; + end; + + PreviousCloseState:=ReceiveClose; + if PreviousCloseState in [csSent,csClosed] then + Result:=False + else if woCloseExplicit in Options then + Result:=True + else + begin + Result:=False; + { Commit/send the close reply before entering application code. } + if aFrame.Reason=CLOSE_NO_STATUS_RCVD then + CloseData:=Nil else - if not (woPongExplicit in Options) then + begin + SetLength(CloseData,2); + CloseData[0]:=(aFrame.Reason and $FF00) shr 8; + CloseData[1]:=aFrame.Reason and $FF; + end; + if not WriteControlFrame(ftClose,CloseData, + WriteErrorMessage) then + begin + ErrorData:=TEncoding.UTF8.GetBytes( + UnicodeString(WriteErrorMessage)); + DispatchEvent(ftClose,nil,ErrorData); + Exit; + end; + end; + { Explicit-close users now receive the close event needed to respond. } + DispatchEvent(ftClose,aFrame,aFrame.Payload.Data); + end; + + ftBinary,ftText: + begin + if FInitialOpcode in [ftText,ftBinary] then + begin + ProtocolError(CLOSE_PROTOCOL_ERROR); + Exit; + end; + if not MessageSizeAllowed(0,Length(aFrame.Payload.Data)) then + begin + ProtocolError(CLOSE_MESSAGE_TOO_BIG); + Exit; + end; + + FInitialOpcode:=aFrame.FrameType; + FMessageContent:=aFrame.Payload.Data; + if aFrame.FinalFrame then begin - Send(ftPong,aFrame.Payload.Data); - DispatchEvent(ftPing,aFrame,aFrame.Payload.Data); + InitialType:=FInitialOpcode; + MessageContent:=FMessageContent; + { All parser state is reset before application code is entered. } + FInitialOpcode:=ftContinuation; + FMessageContent:=[]; + if (InitialType=ftText) and + (not IsValidUTF8(MessageContent)) then + begin + ProtocolError(CLOSE_INVALID_FRAME_PAYLOAD_DATA); + Exit; + end; + DispatchEvent(InitialType,aFrame,MessageContent); end; - end; - ftClose: - begin - // If our side sent the initial close, this is the reply, and we must disconnect (Result=false). - Result:=FCloseState=csNone; - if Result then - begin - if (aFrame.Payload.DataLength>123) then - begin - ProtocolError(CLOSE_PROTOCOL_ERROR); - exit; - end; - - if not (woCloseExplicit in Options) then - begin - if (aFrame.ReasonCLOSE_TLS_HANDSHAKE) and (aFrame.Reason<3000)) then - begin - ProtocolError(CLOSE_PROTOCOL_ERROR); - exit; - end; - if IsValidUTF8(aFrame.Payload.Data) then - begin - DispatchEvent(ftClose,aFrame,aFrame.Payload.Data); - Close('', aFrame.Reason); // Will update state - UpdateCloseState; - Result:=False; // We can disconnect. - end - else - ProtocolError(CLOSE_PROTOCOL_ERROR); - - end - else - UpdateCloseState - end - else - UpdateCloseState; - end; - ftBinary,ftText: - begin - if FInitialOpcode in [ftText, ftBinary] then - begin - ProtocolError(CLOSE_PROTOCOL_ERROR); - Exit; - end; - FInitialOpcode:=aFrame.FrameType; - FMessageContent:=aFrame.Payload.Data; - if aFrame.FinalFrame then - begin - if aFrame.FrameType = ftText then - if IsValidUTF8(aFrame.Payload.Data) then - DispatchEvent(FInitialOpcode,aFrame,aFrame.Payload.Data) - else - ProtocolError(CLOSE_INVALID_FRAME_PAYLOAD_DATA) - else - DispatchEvent(FInitialOpcode,aFrame,aFrame.Payload.Data); - - FInitialOpcode:=ftContinuation; - end; - end; + end; else ; // avoid Compiler warning End; @@ -1476,7 +2079,8 @@ function TWSConnection.HandleIncoming(aFrame: TWSFrame) : Boolean; function TWSConnection.IsValidUTF8(aValue: TBytes): boolean; var - i, len, n, j: integer; + i, len: SizeInt; + n, j: Integer; c: ^byte; begin Result := true; @@ -1585,12 +2189,15 @@ procedure TWSConnection.Close(aMessage: UTF8String); procedure TWSConnection.Close(aMessage: UTF8String; aReason: word); var aData: TBytes; - aSize: Integer; + aSize: SizeInt; begin aData := []; // first two bytes is reason of close RFC 6455 section-5.5.1 aData := TEncoding.UTF8.GetAnsiBytes(aMessage); aSize := Length(aData); + if aSize>123 then + raise EWSProtocolError.Create(SErrControlPayloadTooLarge, + CLOSE_PROTOCOL_ERROR); SetLength(aData, aSize + 2); if aSize > 0 then move(aData[0], aData[2], aSize); @@ -1605,41 +2212,107 @@ procedure TWSConnection.Disconnect; end; procedure TWSConnection.Close(aData: TBytes); +var + CloseReason: Word; + ReasonData: TBytes; begin + if (Length(aData)=1) or (Length(aData)>125) then + raise EWSProtocolError.Create(SErrInvalidPayloadLength, + CLOSE_PROTOCOL_ERROR); + if Length(aData)>=2 then + begin + CloseReason:=NToHs(aData.ToWord(0)); + if not WSValidCloseCode(CloseReason) then + raise EWSProtocolError.Create(SErrInvalidPayloadLength, + CLOSE_PROTOCOL_ERROR); + if Length(aData)>2 then + begin + ReasonData:=Copy(aData,2,Length(aData)-2); + if not IsValidUTF8(ReasonData) then + raise EWSProtocolError.Create(SErrInvalidPayloadLength, + CLOSE_INVALID_FRAME_PAYLOAD_DATA); + end; + end; Send(ftClose,aData); end; -procedure TWSConnection.Send(aFrame: TWSFrame); - -Var - Data : TBytes; - Res: Integer; - ErrMsg: UTF8String; - +function TWSConnection.WriteFrame(aFrame: TWSFrame; + out aErrorMessage: String): Boolean; +var + Data: TBytes; + CurrentTransport: IWSTransport; begin - if FCloseState=csClosed then - Raise EWebSocket.Create(SErrCloseAlreadySent); Data:=aFrame.AsBytes; - Res := Transport.WriteBytes(Data,Length(Data)); - if Res < 0 then - begin - FCloseState:=csClosed; - ErrMsg := Format(SErrWriteReturnedError, [GetLastOSError, SysErrorMessage(GetLastOSError)]); - if woSendErrClosesConn in Options then - begin - SetLength(Data, 0); - Data.Append(TEncoding.UTF8.GetBytes(UnicodeString(ErrMsg))); - DispatchEvent(ftClose, nil, Data); - end - else - Raise EWebSocket.Create(ErrMsg); + CurrentTransport:=Transport; + aErrorMessage:=''; + Result:=True; + EnterCriticalSection(FWriteLock); + try + case FCloseState of + csNone: + ; + csReceived: + if aFrame.FrameType<>ftClose then + Raise EWebSocket.Create(SErrConnectionClosing); + csSent, + csClosed: + Raise EWebSocket.Create(SErrCloseAlreadySent); + end; + try + if not Assigned(CurrentTransport) then + raise EWSWriteError.Create(SErrWriteClosed,0); + { IWSTransport.WriteBuffer is the serialized write-all boundary. + TWSSocketHelper keeps its transport lock across every short write. } + CurrentTransport.WriteBuffer(Data); + except + on E: EWSTransportError do + begin + { Publish terminal state before the caller may dispatch an error. } + FCloseState:=csClosed; + if woSendErrClosesConn in Options then + begin + aErrorMessage:=E.Message; + Result:=False; + end + else + raise; + end; + end; + if Result and (aFrame.FrameType=ftClose) then + begin + if FCloseState=csNone then + FCloseState:=csSent + else if FCloseState=csReceived then + FCloseState:=csClosed; + end; + finally + LeaveCriticalSection(FWriteLock); + end; +end; + +function TWSConnection.WriteControlFrame(aFrameType: TFrameType; + const aData: TBytes; out aErrorMessage: String): Boolean; +var + Frame: TWSFrame; +begin + Frame:=FrameClass.Create(aFrameType,True,aData,OutgoingFrameMask); + try + Result:=WriteFrame(Frame,aErrorMessage); + finally + Frame.Free; end; - if (aFrame.FrameType=ftClose) then +end; + +procedure TWSConnection.Send(aFrame: TWSFrame); +var + ErrorData: TBytes; + ErrorMessage: String; +begin + if not WriteFrame(aFrame,ErrorMessage) then begin - if FCloseState=csNone then - FCloseState:=csSent - else if FCloseState=csReceived then - FCloseState:=csClosed; + ErrorData:=TEncoding.UTF8.GetBytes(UnicodeString(ErrorMessage)); + { This callback is deliberately the final operation on Self. } + DispatchEvent(ftClose,nil,ErrorData); end; end; @@ -1652,12 +2325,31 @@ function TWSConnection.DoReadMessage: Boolean; Result:=False; If not Transport.CanRead(0) then Exit; - f:=FrameClass.CreateFromStream(Transport); try - if Assigned(F) then - Result:=HandleIncoming(F) - finally - F.Free; + F:=Nil; + try + F:=FrameClass.CreateFromStream(Transport,FMaxFramePayloadSize); + if Assigned(F) then + begin + if ((Self is TWSServerConnection) and (not F.Payload.Masked)) or + ((Self is TWSClientConnection) and F.Payload.Masked) then + raise EWSProtocolError.Create(SErrInvalidFrameMask, + CLOSE_PROTOCOL_ERROR); + Result:=HandleIncoming(F); + end; + finally + F.Free; + end; + except + { Preserve the long-standing CheckIncoming contract: an orderly EOF, + including EOF in an incomplete frame, is a terminal irClose result. } + on E: EWSConnectionClosed do + Result:=False; + on E: EWSProtocolError do + begin + Result:=False; + WriteProtocolClose(E.CloseCode); + end; end; end; @@ -1682,7 +2374,7 @@ constructor TWSClientConnection.Create(aOwner: TComponent; aTransport: TWSClient destructor TWSClientConnection.Destroy; begin FreeAndNil(FTransport); - inherited; + inherited Destroy; end; function TWSClientConnection.GetHandshakeCompleted: Boolean; @@ -1790,20 +2482,36 @@ procedure TWSServerConnection.PerformHandshake; Headers : TStrings; aResource,Status,aLine : String; HSR : TWSHandShakeRequest; + HeaderBytes, + LineBytes : SizeInt; + HeaderLines : Integer; begin Status:=Transport.ReadLn; + if Length(Status)>WSMaxServerHandshakeLineBytes then + Raise EWSHandShake.Create(SErrHandshakeLineTooLong); + HeaderBytes:=Length(Status)+2; + HeaderLines:=1; aResource:=ExtractWord(2,Status,[' ']); HSR:=Nil; Headers:=TStringList.Create; try Headers.NameValueSeparator:=':'; - aLine:=Transport.ReadLn; - While aLine<>'' do + repeat begin - Headers.Add(aLine); aLine:=Transport.ReadLn; + LineBytes:=Length(aLine); + if LineBytes>WSMaxServerHandshakeLineBytes then + Raise EWSHandShake.Create(SErrHandshakeLineTooLong); + Inc(HeaderLines); + if (HeaderLines>WSMaxServerHandshakeLines) or + (HeaderBytes+LineBytes+2>WSMaxServerHandshakeBytes) then + Raise EWSHandShake.Create(SErrHandshakeHeadersTooLarge); + Inc(HeaderBytes,LineBytes+2); + if aLine<>'' then + Headers.Add(aLine); end; + until aLine=''; HSR:=TWSHandShakeRequest.Create(aResource,Headers); FHandshakeResponseSent:=DoHandshake(HSR); finally @@ -1863,7 +2571,7 @@ function TWSServerConnection.DoHandshake(const aRequest : TWSHandShakeRequest) : Reply:=Reply+aLine+#13#10; Reply:=Reply+#13#10; B:=TEncoding.UTF8.GetAnsiBytes(Reply); - Transport.WriteBytes(B,Length(B)); + Transport.WriteBuffer(B); Result:=True; FHandshakeResponseSent:=True; except diff --git a/packages/fcl-web/src/websocket/fpwebsocketclient.pp b/packages/fcl-web/src/websocket/fpwebsocketclient.pp index e3b36dc2..6e302c6e 100644 --- a/packages/fcl-web/src/websocket/fpwebsocketclient.pp +++ b/packages/fcl-web/src/websocket/fpwebsocketclient.pp @@ -40,12 +40,22 @@ interface TWSMessagePump = Class (TComponent) private FInterval:Integer; + FCore: TObject; + FEntries: TThreadList; FList: TThreadList; + FRegistryLock: TRTLCriticalSection; FReads: TSocketStreamArray; FExceptions : TSocketStreamArray; FOnError: TWSErrorEvent; + function FindEntry(aConnection: TWSClientConnection): TObject; + procedure RemoveEntry(aEntry: TObject); + procedure StartEntry(aEntry: TObject; aRunGeneration: LongInt); + procedure EntryEnded(aEntry: TObject; aError: Exception); + procedure ReportError(aError: Exception); + procedure ClearEntries; procedure SetInterval(AValue: Integer); Protected + Procedure InterruptConnections; function WaitForData: Boolean; Function CheckConnections : Boolean; virtual; Procedure ReadConnections; @@ -53,6 +63,8 @@ interface Public Constructor Create(aOwner : TComponent); override; Destructor Destroy; override; + // Register a connection created by TCustomWebsocketClient. Its internal + // session supplies the lifetime lease required by the threaded reader. Procedure AddClient(aConnection : TWSClientConnection); Procedure RemoveClient(aConnection : TWSClientConnection); Procedure Execute; virtual; abstract; @@ -66,16 +78,23 @@ interface TWSThreadMessagePump = Class(TWSMessagePump) Private FThread : TThread; - Procedure ThreadTerminated(Sender : TObject); + FLifecycleLock : TRTLCriticalSection; + function TryFinalize(aTimeoutMs: Integer; + aInterruptReaders: Boolean): Boolean; + procedure RequestStop; Protected Type TMessageDriverThread = Class(TThread) Public FPump : TWSThreadMessagePump; - Constructor Create(aPump : TWSThreadMessagePump; aTerminate : TNotifyEvent); + FRunGeneration : LongInt; + Constructor Create(aPump : TWSThreadMessagePump; + aTerminate : TNotifyEvent); Procedure Execute;override; End; Public + Constructor Create(aOwner : TComponent); override; + Destructor Destroy; override; Procedure Execute; override; Procedure Terminate; override; End; @@ -85,10 +104,14 @@ TCustomWebsocketClient = class; { TWebSocketClientConnection } TWebSocketClientConnection = class(TWSClientConnection) + private + FClientSession: TObject; + procedure SetClientSession(aSession: TObject); protected Procedure DoDisconnect; override; function GetClient: TCustomWebsocketClient; virtual; Public + procedure Send(aFrame: TWSFrame); overload; override; Property WebsocketClient : TCustomWebsocketClient Read GetClient; end; @@ -105,6 +128,8 @@ TWebSocketClientConnection = class(TWSClientConnection) FResource: string; FConnectTimeout: Integer; FOptions: TWSOptions; + FMaxFramePayloadSize: QWord; + FMaxMessagePayloadSize: QWord; FSocket : TInetSocket; FTransport : TWSClientTransport; FCheckTimeOut: Integer; @@ -119,6 +144,28 @@ TWebSocketClientConnection = class(TWSClientConnection) FOnControl: TWSControlEvent; FOnDisconnect: TNotifyEvent; FOnConnect: TNotifyEvent; + FStateLock: TRTLCriticalSection; + FDestroying: Boolean; + FNextGeneration: QWord; + FConnectingGeneration: QWord; + FOwnerGate: TObject; + FSession: TObject; + function GetActive: Boolean; + function GetConnection: TWebSocketClientConnection; + function AcquireCurrentSession: TObject; + function AcquireHandshakeSession: TObject; + function IsCurrentSession(aSession: TObject): Boolean; + procedure DisconnectSession(aSession: TObject; SendClose: Boolean; + aEventSender: TObject); + function DetachCurrentSession(aSession: TObject): Boolean; + procedure SessionMessageReceived(aSession: TObject; + const aMessage: TWSMessage); + procedure SessionControlReceived(aSession: TObject; aEventSender: TObject; + aType: TFrameType; const aData: TBytes); + procedure SessionDisconnected(aSession: TObject; aEventSender: TObject; + aPump: TWSMessagePump); + procedure SessionConnected(aSession: TObject); + procedure ReportCallbackError(aPump: TWSMessagePump; E: Exception); procedure FreeConnectionObjects; procedure SetActive(const Value: Boolean); procedure SetHostName(const Value: String); @@ -129,12 +176,16 @@ TWebSocketClientConnection = class(TWSClientConnection) procedure SetResource(const Value: string); procedure SetCheckTimeOut(const Value: Integer); procedure SetOptions(const Value: TWSOptions); + procedure SetMaxFramePayloadSize(const Value: QWord); + procedure SetMaxMessagePayloadSize(const Value: QWord); procedure SetAutoCheckMessages(const Value: Boolean); procedure SendHeaders(aHeaders: TStrings); procedure ConnectionDisconnected(Sender: TObject); Protected Procedure CheckInactive; Procedure Loaded; override; + Procedure Notification(aComponent : TComponent; + Operation : TOperation); override; function CreateClientConnection(aTransport : TWSClientTransport): TWebSocketClientConnection; virtual; procedure MessageReceived(Sender: TObject; const aMessage : TWSMessage); Procedure ControlReceived(Sender: TObject; aType : TFrameType; const aData: TBytes);virtual; @@ -146,8 +197,9 @@ TWebSocketClientConnection = class(TWSClientConnection) Function DoHandShake: Boolean; Property Transport: TWSClientTransport Read FTransport; Public - Property Connection: TWebSocketClientConnection Read FConnection; + Property Connection: TWebSocketClientConnection Read GetConnection; Public + Constructor Create(aOwner : TComponent); override; Destructor Destroy; override; // Check for incoming messages Function CheckIncoming : TIncomingResult; @@ -165,7 +217,7 @@ TWebSocketClientConnection = class(TWSClientConnection) Procedure SendMessage(Const aMessage : String); Public // Connect/Disconnect - Property Active : Boolean Read FActive Write SetActive; + Property Active : Boolean Read GetActive Write SetActive; // Check for message timeout Property CheckTimeOut : Integer Read FCheckTimeOut Write SetCheckTimeOut; // Timeout for connect @@ -176,6 +228,12 @@ TWebSocketClientConnection = class(TWSClientConnection) Property MessagePump : TWSMessagePump Read FMessagePump Write SetMessagePump; // Options Property Options : TWSOptions Read FOptions Write SetOptions; + // Maximum accepted frame payload, 0 means unlimited. + Property MaxFramePayloadSize : QWord Read FMaxFramePayloadSize + Write SetMaxFramePayloadSize; + // Maximum accepted reassembled message payload, 0 means unlimited. + Property MaxMessagePayloadSize : QWord Read FMaxMessagePayloadSize + Write SetMaxMessagePayloadSize; // Mask to use for outgoing frames Property OutGoingFrameMask : Integer Read FOutGoingFrameMask Write FOutGoingFrameMask; // Port to connect to @@ -207,6 +265,8 @@ TWebSocketClientConnection = class(TWSClientConnection) Property ConnectTimeout; Property MessagePump; Property Options; + Property MaxFramePayloadSize; + Property MaxMessagePayloadSize; Property Resource; Property UseSSL; Property OnSendHandShake; @@ -226,346 +286,2153 @@ implementation uses sha1; {$ENDIF FPC_DOTTEDUNITS} -{ TWebSocketClientConnection } +{ Internal design overview + + The public client/component model is retained, but connections are managed + internally as generations: + + * TWSClientSession represents one successful or in-progress Connect. It + owns the connection, transport, socket and handshake objects. The client, + pump registration, reader and temporary API operations hold counted + leases, so an old generation can finish without touching a reconnect. + * TWSMessagePumpEntry is the stable registry identity. Removal is by entry + or connection identity, never by an index retained across a callback. + * TWSPumpCore outlives the component while workers drain. Workers acquire a + short owner admission before calling pump code, so pump destruction cannot + race EntryEnded or OnError. + * TWSClientOwnerGate protects all callbacks into TCustomWebsocketClient. + Destruction closes admission and waits for admitted callbacks, servicing + CheckSynchronize on the main thread while it waits. + + The threaded pump uses one blocking reader per registered session. Its + driver still calls the protected CheckConnections/ReadConnections pipeline, + but ReadConnections schedules readers rather than performing a competing + read. This prevents one incomplete frame from starving healthy siblings. + + Terminal notification first acquires callback admission, then claims the + disconnect once, marks the exact session closing, updates client-visible + state, removes its registry entry with durable read cancellation, and + finally invokes OnDisconnect. Application callbacks run without the + registry, state or transport write locks held. + + Disconnect and reconnect are supported from message, control and disconnect + callbacks, including methods reached through TThread.Synchronize. Directly + freeing a client or pump from its own callback is rejected: Pascal cannot + safely return from that callback through an already destroyed object. + + TWSMessagePump.AddClient accepts the public base parameter for compatibility + but requires a TWebSocketClientConnection carrying this unit's managed + session. A raw TWSClientConnection has no lifetime lease and therefore + cannot be read safely while an unrelated owner may destroy it. + + See README.md in this directory for the full invariants, shutdown algorithm, + compatibility decisions and regression-test recipe. } + +Const + WSClientSessionOpen = 0; + WSClientSessionClosing = 1; + WSClientSessionClosed = 2; + WSPumpEntryRegistered = 0; + WSPumpEntryRemoved = 1; + WSMaxHandshakeHeaderLineBytes = 8192; + WSMaxHandshakeHeaderLines = 256; + WSMaxHandshakeHeaderBytes = 65536; + +Resourcestring + SErrClientSessionClosing = 'WebSocket client session is closing'; + SErrFreeFromCallback = 'A websocket client cannot be freed from its own callback'; + SErrHandshakeHeaderLineTooLong = 'WebSocket handshake response header line is too long'; + SErrHandshakeHeadersTooLarge = 'WebSocket handshake response headers are too large'; -procedure TWebSocketClientConnection.DoDisconnect; -begin - If Assigned(WebSocketClient) then - WebSocketClient.ConnectionDisconnected(Self); -end; +Type + TWSPumpCore = Class; + TWSMessagePumpEntry = Class; + TWSClientOwnerGate = Class; + PWSOwnerGateToken = ^TWSOwnerGateToken; + TWSOwnerGateToken = Record + Gate : TWSClientOwnerGate; + Previous : PWSOwnerGateToken; + end; -function TWebSocketClientConnection.GetClient: TCustomWebsocketClient; + { TWSClientOwnerGate -begin - Result:=Owner as TCustomWebsocketClient; -end; + This object, rather than a connection method pointer, protects callbacks + into the component. It is shared by every generation of a client and may + therefore outlive the component itself. } + + TWSClientOwnerGate = Class + Private + FReferenceCount : LongInt; + FLock : TRTLCriticalSection; + FNoCallbacks : PRTLEvent; + FClient : TCustomWebsocketClient; + FAdmissionOpen : Boolean; + FCallbackCount : LongInt; + Public + Constructor Create(aClient : TCustomWebsocketClient); + Destructor Destroy; override; + Procedure AddReference; + Procedure ReleaseReference; + Function TryEnter(out aClient : TCustomWebsocketClient; + out aToken : TWSOwnerGateToken) : Boolean; + Procedure Leave(var aToken : TWSOwnerGateToken); + Procedure CloseAdmission; + Procedure WaitForQuiescence; + Function HasCallbacks : Boolean; + Function IsCurrentCallback : Boolean; + end; + { TWSClientSession -{ TCustomWebsocketClient } + One instance represents exactly one Connect generation. It owns the + concrete connection (and through it the transport and socket). } -procedure TCustomWebsocketClient.CheckInactive; + TWSClientSession = Class + Private + FReferenceCount : LongInt; + FState : LongInt; + FDisconnectClaimed : LongInt; + FGeneration : QWord; + FOwnerGate : TWSClientOwnerGate; + FConnection : TWebSocketClientConnection; + FHandshakeRequest : TWSHandShakeRequest; + FHandshakeResponse : TWSHandShakeResponse; + Procedure MessageReceived(Sender : TObject; const aMessage : TWSMessage); + Procedure ControlReceived(Sender : TObject; aType : TFrameType; + const aData : TBytes); + Public + Constructor Create(aOwnerGate : TWSClientOwnerGate; + aGeneration : QWord); + Destructor Destroy; override; + Procedure AddReference; + Procedure ReleaseReference; + Procedure AttachConnection(aConnection : TWebSocketClientConnection); + Procedure SetHandshakeRequest(aRequest : TWSHandShakeRequest); + Procedure SetHandshakeResponse(aResponse : TWSHandShakeResponse); + Function TakeHandshakeResponse : TWSHandShakeResponse; + Function IsOpen : Boolean; + Function BeginClosing : Boolean; + Procedure CancelReads; + Procedure InterruptRead; + Procedure ConnectionRequestedDisconnect(aEventSender : TObject); + Procedure NotifyDisconnected(aPump : TWSMessagePump; + aEventSender : TObject); + Procedure NotifyConnected; + Property Connection : TWebSocketClientConnection Read FConnection; + Property Generation : QWord Read FGeneration; + Property HandshakeRequest : TWSHandShakeRequest Read FHandshakeRequest; + end; + + { Stable pump registration. Registry, worker and temporary users each own + a reference. The entry owns one session reference. } + + TWSMessagePumpEntry = Class + Private + FReferenceCount : LongInt; + FState : LongInt; + FWorkerRunning : LongInt; + FSession : TWSClientSession; + FConnection : TWSClientConnection; + Public + Constructor Create(aSession : TWSClientSession; + aConnection : TWSClientConnection); + Destructor Destroy; override; + Procedure AddReference; + Procedure ReleaseReference; + Function MarkRemoved : Boolean; + Function IsRegistered : Boolean; + Function TryStartWorker : Boolean; + Procedure WorkerStopped; + Function WorkerRunning : Boolean; + Procedure CancelReads; + Procedure InterruptRead; + Property Connection : TWSClientConnection Read FConnection; + Property Session : TWSClientSession Read FSession; + end; + + { Refcounted state used by reader workers. A worker never retains the + component pointer directly; it obtains a short owner admission here. } + + TWSPumpCore = Class + Private + FReferenceCount : LongInt; + FLock : TRTLCriticalSection; + FNoWorkers : PRTLEvent; + FNoOwnerUsers : PRTLEvent; + FPump : TWSMessagePump; + FOwnerOpen : Boolean; + FOwnerUsers : LongInt; + FWorkerCount : LongInt; + FRunning : Boolean; + FRunGeneration : LongInt; + FInterval : Integer; + Public + Constructor Create(aPump : TWSMessagePump; aInterval : Integer); + Destructor Destroy; override; + Procedure AddReference; + Procedure ReleaseReference; + Function BeginRun(out aGeneration : LongInt) : Boolean; + Procedure RequestStop; + Function IsRunning(aGeneration : LongInt) : Boolean; + Function CurrentGeneration : LongInt; + Procedure WorkerStarting; + Procedure WorkerDone; + Function WorkerCount : LongInt; + Function WaitWorkers(aTimeoutMs : Integer) : Boolean; + Function TryEnterPump(out aPump : TWSMessagePump) : Boolean; + Procedure LeavePump; + Procedure CloseOwner; + Procedure WaitOwnerUsers; + Procedure SetInterval(aValue : Integer); + Function GetInterval : Integer; + end; + + TWSClientReaderThread = Class(TThread) + Private + FCore : TWSPumpCore; + FEntry : TWSMessagePumpEntry; + FRunGeneration : LongInt; + Procedure FinishWithError(aError : Exception); + Public + Constructor Create(aCore : TWSPumpCore; aEntry : TWSMessagePumpEntry; + aRunGeneration : LongInt); + Procedure AbandonBeforeStart; + Procedure Execute; override; + end; + +ThreadVar + CurrentWSPumpCore : TWSPumpCore; + CurrentWSOwnerGateToken : PWSOwnerGateToken; + CurrentWSHandshakeClient : TCustomWebsocketClient; + CurrentWSHandshakeSession : TWSClientSession; + +procedure ReleaseClientSession(aSession: TWSClientSession; + aPump: TWSMessagePump); begin - If Active then - Raise EWebSocketClient.Create(SErrConnectionActive); + if not Assigned(aSession) then + Exit; + try + aSession.ReleaseReference; + except + on E : Exception do + if Assigned(aPump) then + aPump.ReportError(E); + end; end; -Function TCustomWebsocketClient.CheckIncoming : TIncomingResult; +{ TWSClientOwnerGate } +constructor TWSClientOwnerGate.Create(aClient: TCustomWebsocketClient); begin - If Not Active then - Raise EWebSocketClient.Create(SErrConnectionInActive); - if Not Connection.HandshakeCompleted then - Raise EWebSocketClient.Create(SErrHandshakeInComplete); - Result:=Connection.CheckIncoming(CheckTimeout); - if (Result=irClose) then - begin - Disconnect(False); - end; + inherited Create; + FReferenceCount:=1; + InitCriticalSection(FLock); + FNoCallbacks:=RTLEventCreate; + RTLEventSetEvent(FNoCallbacks); + FClient:=aClient; + FAdmissionOpen:=True; end; -procedure TCustomWebsocketClient.ControlReceived(Sender: TObject; aType : TFrameType; const aData: TBytes); +destructor TWSClientOwnerGate.Destroy; begin - If Assigned(FOnControl) then - FOnControl(Sender, aType, aData); + RTLEventDestroy(FNoCallbacks); + DoneCriticalSection(FLock); + inherited Destroy; end; -function TCustomWebsocketClient.CreateClientConnection(aTransport: TWSClientTRansport): TWebsocketClientConnection; - +procedure TWSClientOwnerGate.AddReference; begin - Result:=TWebSocketClientConnection.Create(Self,aTransport,FOptions); + InterlockedIncrement(FReferenceCount); end; -procedure TCustomWebsocketClient.ConnectionDisconnected(Sender : TObject); - +procedure TWSClientOwnerGate.ReleaseReference; begin - FActive:=False; - If Assigned(MessagePump) then - MessagePump.RemoveClient(FConnection); - If Assigned(OnDisconnect) then - OnDisconnect(FConnection); - // We cannot free the connection here, because it still needs to call it's own OnDisconnect. + if InterlockedDecrement(FReferenceCount)=0 then + Free; end; -procedure TCustomWebsocketClient.Connect; -var - SSLHandler: TSSLSocketHandler; +function TWSClientOwnerGate.TryEnter( + out aClient: TCustomWebsocketClient; + out aToken: TWSOwnerGateToken): Boolean; begin - If Active then - Exit; - // Safety: Free any dangling objects before recreating - FreeConnectionObjects; - SSLHandler := nil; - if UseSSL then - begin - SSLHandler := TSSLSocketHandler.GetDefaultHandler; - SSLHandler.VerifyPeerCert := False; + aClient:=Nil; + aToken.Gate:=Nil; + aToken.Previous:=Nil; + EnterCriticalSection(FLock); + try + Result:=FAdmissionOpen and Assigned(FClient); + if Result then + begin + if FCallbackCount=0 then + RTLEventResetEvent(FNoCallbacks); + Inc(FCallbackCount); + aClient:=FClient; + end; + finally + LeaveCriticalSection(FLock); end; - FSocket:=TInetSocket.Create(HostName,Port,ConnectTimeout, SSLHandler); - FTransport:=TWSClientTransport.Create(FSocket); - FConnection:=CreateClientConnection(FTransport); - FConnection.OnMessageReceived:=@MessageReceived; - FConnection.OnControl:=@ControlReceived; - // RFC states we MUST use a mask. - if OutGoingFrameMask=0 then - OutGoingFrameMask:=1+Random(MaxInt-1); - FConnection.OutgoingFrameMask:=Self.OutGoingFrameMask; - if UseSSL then - FSocket.Connect; - FActive:=True; - if not DoHandShake then - Disconnect(False) - else + if Result then begin - If Assigned(MessagePump) then - MessagePump.AddClient(FConnection); - if Assigned(OnConnect) then - OnConnect(Self); + aToken.Gate:=Self; + aToken.Previous:=CurrentWSOwnerGateToken; + CurrentWSOwnerGateToken:=@aToken; end; end; - -destructor TCustomWebsocketClient.Destroy; -begin - DisConnect(False); - FreeAndNil(FHandShake); - FreeAndNil(FHandshakeResponse); - FreeConnectionObjects; - Inherited; -end; - - -Function TCustomWebsocketClient.CreateHandShakeRequest : TWSHandShakeRequest; - +procedure TWSClientOwnerGate.Leave(var aToken: TWSOwnerGateToken); begin - Result:=TWSHandShakeRequest.Create('',Nil); + EnterCriticalSection(FLock); + try + Dec(FCallbackCount); + if FCallbackCount=0 then + RTLEventSetEvent(FNoCallbacks); + finally + LeaveCriticalSection(FLock); + end; + if CurrentWSOwnerGateToken=@aToken then + CurrentWSOwnerGateToken:=aToken.Previous; + aToken.Gate:=Nil; + aToken.Previous:=Nil; end; -procedure TCustomWebsocketClient.SendData(aBytes: TBytes); - +procedure TWSClientOwnerGate.CloseAdmission; begin - Connection.Send(aBytes); + EnterCriticalSection(FLock); + try + FAdmissionOpen:=False; + FClient:=Nil; + finally + LeaveCriticalSection(FLock); + end; end; -procedure TCustomWebsocketClient.SendHeaders(aHeaders : TStrings); - -Var - S : String; - B : TBytes; - +procedure TWSClientOwnerGate.WaitForQuiescence; +var + Pending : LongInt; begin - for S in AHeaders do - begin - B:=TEncoding.UTF8.GetAnsiBytes(S+#13#10); - Connection.Transport.WriteBytes(B,Length(B)); + repeat + EnterCriticalSection(FLock); + try + Pending:=FCallbackCount; + finally + LeaveCriticalSection(FLock); end; - B:=TEncoding.UTF8.GetAnsiBytes(#13#10); - Connection.Transport.WriteBytes(B,Length(B)); + if Pending=0 then + Exit; + if TThread.CurrentThread.ThreadID=MainThreadID then + CheckSynchronize(0); + RTLEventWaitFor(FNoCallbacks,1); + until False; end; -procedure TCustomWebsocketClient.SendHandShakeRequest; - -Var - aRequest : TWSHandShakeRequest; - aHeaders : TStrings; +function TWSClientOwnerGate.HasCallbacks: Boolean; begin - aHeaders:=Nil; - FreeAndNil(FHandShake); - aRequest:=CreateHandShakeRequest; + EnterCriticalSection(FLock); try - aRequest.Host:=HostName; - aRequest.Port:=Port; - aRequest.Resource:=Resource; - aHeaders:=TStringList.Create; - aHeaders.NameValueSeparator:=':'; - aRequest.ToStrings(aHeaders); - if Assigned(FOnSendHandshake) then - FOnSendHandshake(self,aHeaders); - // Do not use FClient.WriteHeader, it messes up the strings ! - SendHeaders(aHeaders); - FHandShake:=aRequest; + Result:=FCallbackCount<>0; finally - aHeaders.Free; - if FhandShake<>aRequest then - aRequest.Free; + LeaveCriticalSection(FLock); end; end; -procedure TCustomWebsocketClient.SendMessage(const aMessage: String); +function TWSClientOwnerGate.IsCurrentCallback: Boolean; +var + Token : PWSOwnerGateToken; begin - Connection.Send(aMessage); + Token:=CurrentWSOwnerGateToken; + while Assigned(Token) and (Token^.Gate<>Self) do + Token:=Token^.Previous; + Result:=Assigned(Token); end; -Function TCustomWebsocketClient.CreateHandshakeResponse(aHeaders : TStrings) : TWSHandShakeResponse; +{ TWSClientSession } +constructor TWSClientSession.Create(aOwnerGate: TWSClientOwnerGate; + aGeneration: QWord); begin - Result:=TWSHandShakeResponse.Create('',aHeaders); + inherited Create; + FReferenceCount:=1; + FState:=WSClientSessionOpen; + FGeneration:=aGeneration; + FOwnerGate:=aOwnerGate; + FOwnerGate.AddReference; end; -Function TCustomWebsocketClient.CheckHandShakeResponse(aHeaders : TStrings) : Boolean; +destructor TWSClientSession.Destroy; +var + ConnectionToFree : TWebSocketClientConnection; +begin + InterlockedExchange(FState,WSClientSessionClosed); + ConnectionToFree:=FConnection; + FConnection:=Nil; + if Assigned(ConnectionToFree) then + begin + ConnectionToFree.SetClientSession(Nil); + ConnectionToFree.OnMessageReceived:=Nil; + ConnectionToFree.OnControl:=Nil; + ConnectionToFree.HandshakeResponse:=Nil; + end; + try + try + FreeAndNil(FHandshakeResponse); + finally + try + FreeAndNil(FHandshakeRequest); + finally + if Assigned(ConnectionToFree) then + begin + ConnectionToFree.Free; + end; + end; + end; + finally + try + FOwnerGate.ReleaseReference; + finally + inherited Destroy; + end; + end; +end; -Var - K : String; - {%H-}hash : TSHA1Digest; - B : TBytes; +procedure TWSClientSession.SetHandshakeRequest( + aRequest: TWSHandShakeRequest); +begin + if FHandshakeRequest=aRequest then + Exit; + FreeAndNil(FHandshakeRequest); + FHandshakeRequest:=aRequest; +end; +procedure TWSClientSession.SetHandshakeResponse( + aResponse: TWSHandShakeResponse); begin - B:=[]; + if FHandshakeResponse=aResponse then + Exit; FreeAndNil(FHandshakeResponse); - FHandshakeResponse:=CreateHandshakeResponse(aHeaders); - k := Trim(FHandshake.Key) + SSecWebSocketGUID; - hash:=SHA1String(k); - SetLength(B,SizeOf(hash)); - Move(hash[0],B[0],SizeOf(hash)); - k:=EncodeBytesBase64(B); - Result:=SameText(K,FHandshakeResponse.Accept) - and SameText(FHandshakeResponse.Upgrade,'websocket'); + FHandshakeResponse:=aResponse; end; -Function TCustomWebsocketClient.ReadHandShakeResponse : Boolean; +function TWSClientSession.TakeHandshakeResponse: TWSHandShakeResponse; +begin + Result:=FHandshakeResponse; + FHandshakeResponse:=Nil; +end; -Var - S : String; - aHeaders : TStrings; +procedure TWSClientSession.AddReference; +begin + InterlockedIncrement(FReferenceCount); +end; +procedure TWSClientSession.ReleaseReference; begin - Result:=False; - aHeaders:=TStringList.Create; - Try - aHeaders.NameValueSeparator:=':'; - Repeat - S:=Connection.Transport.ReadLn; - aHeaders.Add(S); - Until (S=''); - Result:=CheckHandShakeResponse(aHeaders); - if Result and Assigned(FOnHandshakeResponse) then - FOnHandshakeResponse(Self,FHandShakeResponse,Result); - if Result then - FConnection.HandshakeResponse:=FHandShakeResponse - Finally - aHeaders.Free; - End; + if InterlockedDecrement(FReferenceCount)=0 then + Free; end; -Function TCustomWebsocketClient.DoHandShake : Boolean; +procedure TWSClientSession.AttachConnection( + aConnection: TWebSocketClientConnection); +begin + if Assigned(FConnection) then + Raise EWebSocketClient.Create('A websocket session already has a connection'); + FConnection:=aConnection; + FConnection.SetClientSession(Self); + FConnection.OnMessageReceived:=@MessageReceived; + FConnection.OnControl:=@ControlReceived; +end; +function TWSClientSession.IsOpen: Boolean; begin - SendHandShakeRequest; - Result:=ReadHandShakeResponse; + Result:=InterlockedCompareExchange(FState,WSClientSessionOpen, + WSClientSessionOpen)=WSClientSessionOpen; end; -procedure TCustomWebsocketClient.Loaded; +function TWSClientSession.BeginClosing: Boolean; begin - inherited; - if FLoadActive then - Connect; + Result:=InterlockedCompareExchange(FState,WSClientSessionClosing, + WSClientSessionOpen)=WSClientSessionOpen; end; -procedure TCustomWebsocketClient.MessageReceived(Sender: TObject; const aMessage : TWSMessage) ; +procedure TWSClientSession.CancelReads; begin - if Assigned(OnMessageReceived) and (TWSClientConnection(Sender).HandshakeCompleted) then - OnMessageReceived(Self, AMessage); + if Assigned(FConnection) and Assigned(FConnection.ClientTransport) then + FConnection.ClientTransport.CancelReads; end; -procedure TCustomWebsocketClient.Ping(aMessage: UTF8String); +procedure TWSClientSession.InterruptRead; begin - FConnection.Send(ftPing,TEncoding.UTF8.GetAnsiBytes(aMessage)); + if Assigned(FConnection) and Assigned(FConnection.ClientTransport) then + FConnection.ClientTransport.InterruptRead; end; -procedure TCustomWebsocketClient.Pong(aMessage: UTF8String); +procedure TWSClientSession.MessageReceived(Sender: TObject; + const aMessage: TWSMessage); +var + aClient : TCustomWebsocketClient; + GateToken : TWSOwnerGateToken; begin - FConnection.Send(ftPong,TEncoding.UTF8.GetAnsiBytes(aMessage)); + if not IsOpen then + Exit; + if not FOwnerGate.TryEnter(aClient,GateToken) then + Exit; + try + if IsOpen then + aClient.SessionMessageReceived(Self,aMessage); + finally + FOwnerGate.Leave(GateToken); + end; end; -procedure TCustomWebsocketClient.FreeConnectionObjects; +procedure TWSClientSession.ControlReceived(Sender: TObject; aType: TFrameType; + const aData: TBytes); +var + aClient : TCustomWebsocketClient; + GateToken : TWSOwnerGateToken; +begin + if not IsOpen then + Exit; + if not FOwnerGate.TryEnter(aClient,GateToken) then + Exit; + try + if IsOpen then + aClient.SessionControlReceived(Self,Sender,aType,aData); + finally + FOwnerGate.Leave(GateToken); + end; +end; +procedure TWSClientSession.ConnectionRequestedDisconnect(aEventSender: TObject); begin - FreeAndNil(FConnection); - FTransport:=nil; // FTransport is freed in TWSClientConnection.Destroy - FSocket:=nil; // FSocket is freed in TWSClientTransport.Destroy + AddReference; + try + NotifyDisconnected(Nil,aEventSender); + finally + ReleaseClientSession(Self,Nil); + end; end; -procedure TCustomWebsocketClient.Disconnect(SendClose : boolean = true); +procedure TWSClientSession.NotifyDisconnected(aPump: TWSMessagePump; + aEventSender: TObject); +var + aClient : TCustomWebsocketClient; + GateToken : TWSOwnerGateToken; +begin + { Admission is deliberately acquired before the once-only claim. Closing + a client can then either wait for this admitted notification or claim and + deliver the notification itself; it can never be silently lost between + the claim and a closing callback gate. } + if not FOwnerGate.TryEnter(aClient,GateToken) then + Exit; + try + if InterlockedCompareExchange(FDisconnectClaimed,1,0)=0 then + aClient.SessionDisconnected(Self,aEventSender,aPump); + finally + FOwnerGate.Leave(GateToken); + end; +end; +procedure TWSClientSession.NotifyConnected; +var + aClient : TCustomWebsocketClient; + GateToken : TWSOwnerGateToken; begin - if Not Active then + if not FOwnerGate.TryEnter(aClient,GateToken) then Exit; - if SendClose and (Connection.CloseState <> csClosed) then - Connection.Close(''); - if Assigned(MessagePump) then - MessagePump.RemoveClient(Connection); - If Assigned(OnDisconnect) then - OnDisconnect(Self); - FreeConnectionObjects; - FActive:=False; + try + if IsOpen then + aClient.SessionConnected(Self); + finally + FOwnerGate.Leave(GateToken); + end; end; -procedure TCustomWebsocketClient.SetActive(const Value: Boolean); +{ TWSMessagePumpEntry } + +constructor TWSMessagePumpEntry.Create(aSession: TWSClientSession; + aConnection: TWSClientConnection); begin - FLoadActive := Value; - if (csDesigning in ComponentState) then - exit; - if Value then - Connect - else - Disconnect; + inherited Create; + FReferenceCount:=1; + FState:=WSPumpEntryRegistered; + FSession:=aSession; + FSession.AddReference; + FConnection:=aConnection; end; -procedure TCustomWebsocketClient.SetAutoCheckMessages(const Value: Boolean); +destructor TWSMessagePumpEntry.Destroy; begin - CheckInactive; - FAutoCheckMessages := Value; + try + FSession.ReleaseReference; + finally + inherited Destroy; + end; end; -procedure TCustomWebsocketClient.SetCheckTimeOut(const Value: Integer); +procedure TWSMessagePumpEntry.AddReference; begin - CheckInactive; - FCheckTimeOut := Value; + InterlockedIncrement(FReferenceCount); end; -procedure TCustomWebsocketClient.SetConnectTimeout(const Value: Integer); +procedure TWSMessagePumpEntry.ReleaseReference; begin - CheckInactive; - FConnectTimeout := Value; + if InterlockedDecrement(FReferenceCount)=0 then + Free; end; -procedure TCustomWebsocketClient.SetHostName(const Value: String); +function TWSMessagePumpEntry.MarkRemoved: Boolean; begin - CheckInactive; - FHostName := Value; + Result:=InterlockedCompareExchange(FState,WSPumpEntryRemoved, + WSPumpEntryRegistered)=WSPumpEntryRegistered; end; -procedure TCustomWebsocketClient.SetMessagePump(AValue: TWSMessagePump); +function TWSMessagePumpEntry.IsRegistered: Boolean; begin - if FMessagePump=AValue then Exit; - If Assigned(FMessagePump) then - FMessagePump.RemoveFreeNotification(Self); - FMessagePump:=AValue; - If Assigned(FMessagePump) then - FMessagePump.FreeNotification(Self); + Result:=InterlockedCompareExchange(FState,WSPumpEntryRegistered, + WSPumpEntryRegistered)=WSPumpEntryRegistered; end; -procedure TCustomWebsocketClient.SetOptions(const Value: TWSOptions); +function TWSMessagePumpEntry.TryStartWorker: Boolean; begin - CheckInactive; - FOptions := Value; + Result:=IsRegistered and FSession.IsOpen and + (InterlockedCompareExchange(FWorkerRunning,1,0)=0); + if Result and ((not IsRegistered) or (not FSession.IsOpen)) then + begin + InterlockedExchange(FWorkerRunning,0); + Result:=False; + end; end; -procedure TCustomWebsocketClient.SetPort(const Value: Integer); +procedure TWSMessagePumpEntry.WorkerStopped; begin - CheckInactive; - FPort := Value; + InterlockedExchange(FWorkerRunning,0); end; -procedure TCustomWebsocketClient.SetResource(const Value: string); +function TWSMessagePumpEntry.WorkerRunning: Boolean; begin - CheckInactive; - FResource := Value; + Result:=InterlockedCompareExchange(FWorkerRunning,0,0)<>0; end; -procedure TCustomWebsocketClient.SetUseSSL(const Value: Boolean); +procedure TWSMessagePumpEntry.CancelReads; begin - CheckInactive; - FUseSSL := Value; + FSession.CancelReads; +end; + +procedure TWSMessagePumpEntry.InterruptRead; +begin + FSession.InterruptRead; +end; + +{ TWSPumpCore } + +constructor TWSPumpCore.Create(aPump: TWSMessagePump; aInterval: Integer); +begin + inherited Create; + FReferenceCount:=1; + InitCriticalSection(FLock); + FNoWorkers:=RTLEventCreate; + FNoOwnerUsers:=RTLEventCreate; + RTLEventSetEvent(FNoWorkers); + RTLEventSetEvent(FNoOwnerUsers); + FPump:=aPump; + FOwnerOpen:=True; + FInterval:=aInterval; +end; + +destructor TWSPumpCore.Destroy; +begin + RTLEventDestroy(FNoOwnerUsers); + RTLEventDestroy(FNoWorkers); + DoneCriticalSection(FLock); + inherited Destroy; +end; + +procedure TWSPumpCore.AddReference; +begin + InterlockedIncrement(FReferenceCount); +end; + +procedure TWSPumpCore.ReleaseReference; +begin + if InterlockedDecrement(FReferenceCount)=0 then + Free; +end; + +function TWSPumpCore.BeginRun(out aGeneration: LongInt): Boolean; +begin + EnterCriticalSection(FLock); + try + Result:=(not FRunning) and (FWorkerCount=0) and FOwnerOpen; + if Result then + begin + Inc(FRunGeneration); + if FRunGeneration=0 then + Inc(FRunGeneration); + FRunning:=True; + end; + aGeneration:=FRunGeneration; + finally + LeaveCriticalSection(FLock); + end; +end; + +procedure TWSPumpCore.RequestStop; +begin + EnterCriticalSection(FLock); + try + FRunning:=False; + finally + LeaveCriticalSection(FLock); + end; +end; + +function TWSPumpCore.IsRunning(aGeneration: LongInt): Boolean; +begin + EnterCriticalSection(FLock); + try + Result:=FRunning and (FRunGeneration=aGeneration); + finally + LeaveCriticalSection(FLock); + end; +end; + +function TWSPumpCore.CurrentGeneration: LongInt; +begin + EnterCriticalSection(FLock); + try + Result:=FRunGeneration; + finally + LeaveCriticalSection(FLock); + end; +end; + +procedure TWSPumpCore.WorkerStarting; +begin + EnterCriticalSection(FLock); + try + if FWorkerCount=0 then + RTLEventResetEvent(FNoWorkers); + Inc(FWorkerCount); + finally + LeaveCriticalSection(FLock); + end; +end; + +procedure TWSPumpCore.WorkerDone; +begin + EnterCriticalSection(FLock); + try + Dec(FWorkerCount); + if FWorkerCount=0 then + RTLEventSetEvent(FNoWorkers); + finally + LeaveCriticalSection(FLock); + end; +end; + +function TWSPumpCore.WorkerCount: LongInt; +begin + EnterCriticalSection(FLock); + try + Result:=FWorkerCount; + finally + LeaveCriticalSection(FLock); + end; +end; + +function TWSPumpCore.WaitWorkers(aTimeoutMs: Integer): Boolean; +var + Started : QWord; +begin + Started:=TThread.GetTickCount64; + repeat + Result:=WorkerCount=0; + if Result then + Exit; + if (aTimeoutMs>=0) and + ((TThread.GetTickCount64-Started)>=QWord(aTimeoutMs)) then + Exit(False); + if TThread.CurrentThread.ThreadID=MainThreadID then + CheckSynchronize(0); + RTLEventWaitFor(FNoWorkers,1); + until False; +end; + +function TWSPumpCore.TryEnterPump(out aPump: TWSMessagePump): Boolean; +begin + aPump:=Nil; + EnterCriticalSection(FLock); + try + Result:=FOwnerOpen and Assigned(FPump); + if Result then + begin + if FOwnerUsers=0 then + RTLEventResetEvent(FNoOwnerUsers); + Inc(FOwnerUsers); + aPump:=FPump; + end; + finally + LeaveCriticalSection(FLock); + end; +end; + +procedure TWSPumpCore.LeavePump; +begin + EnterCriticalSection(FLock); + try + Dec(FOwnerUsers); + if FOwnerUsers=0 then + RTLEventSetEvent(FNoOwnerUsers); + finally + LeaveCriticalSection(FLock); + end; +end; + +procedure TWSPumpCore.CloseOwner; +begin + EnterCriticalSection(FLock); + try + FOwnerOpen:=False; + FPump:=Nil; + finally + LeaveCriticalSection(FLock); + end; +end; + +procedure TWSPumpCore.WaitOwnerUsers; +var + Pending : LongInt; +begin + repeat + EnterCriticalSection(FLock); + try + Pending:=FOwnerUsers; + finally + LeaveCriticalSection(FLock); + end; + if Pending=0 then + Exit; + if TThread.CurrentThread.ThreadID=MainThreadID then + CheckSynchronize(0); + RTLEventWaitFor(FNoOwnerUsers,1); + until False; +end; + +procedure TWSPumpCore.SetInterval(aValue: Integer); +begin + EnterCriticalSection(FLock); + try + FInterval:=aValue; + finally + LeaveCriticalSection(FLock); + end; +end; + +function TWSPumpCore.GetInterval: Integer; +begin + EnterCriticalSection(FLock); + try + Result:=FInterval; + finally + LeaveCriticalSection(FLock); + end; +end; + +{ TWSClientReaderThread } + +constructor TWSClientReaderThread.Create(aCore: TWSPumpCore; + aEntry: TWSMessagePumpEntry; aRunGeneration: LongInt); +begin + inherited Create(True); + FCore:=aCore; + FCore.AddReference; + FEntry:=aEntry; + FEntry.AddReference; + FRunGeneration:=aRunGeneration; + FCore.WorkerStarting; + FreeOnTerminate:=True; +end; + +procedure TWSClientReaderThread.AbandonBeforeStart; +var + aPump : TWSMessagePump; + PreviousCore : TWSPumpCore; +begin + PreviousCore:=CurrentWSPumpCore; + CurrentWSPumpCore:=FCore; + try + FEntry.WorkerStopped; + try + try + FEntry.ReleaseReference; + except + on E : Exception do + if FCore.TryEnterPump(aPump) then + try + aPump.ReportError(E); + finally + FCore.LeavePump; + end; + end; + finally + FCore.WorkerDone; + end; + finally + try + FCore.ReleaseReference; + finally + CurrentWSPumpCore:=PreviousCore; + FEntry:=Nil; + FCore:=Nil; + end; + end; +end; + +procedure TWSClientReaderThread.FinishWithError(aError: Exception); +var + aPump : TWSMessagePump; +begin + if FCore.TryEnterPump(aPump) then + try + aPump.EntryEnded(FEntry,aError); + finally + FCore.LeavePump; + end; +end; + +procedure TWSClientReaderThread.Execute; +var + IncomingResult : TIncomingResult; + FinishedConnection : Boolean; + aPump : TWSMessagePump; + PreviousCore : TWSPumpCore; +begin + PreviousCore:=CurrentWSPumpCore; + CurrentWSPumpCore:=FCore; + try + FinishedConnection:=False; + while FCore.IsRunning(FRunGeneration) and FEntry.IsRegistered and + FEntry.Session.IsOpen do + begin + try + IncomingResult:=FEntry.Connection.CheckIncoming(FCore.GetInterval); + if IncomingResult=irClose then + begin + FinishWithError(Nil); + FinishedConnection:=True; + Break; + end; + except + on E : EWSReadInterrupted do + begin + { An interruption observed by the exact reader invalidates a + partially consumed frame. Removal/notification is idempotent. } + FinishWithError(Nil); + FinishedConnection:=True; + Break; + end; + on E : Exception do + begin + FinishWithError(E); + FinishedConnection:=True; + Break; + end; + end; + end; + { Avoid a warning in compilers which do not optimize the loop flag. } + if FinishedConnection then + ; + finally + try + FEntry.WorkerStopped; + try + try + FEntry.ReleaseReference; + except + on E : Exception do + begin + if FCore.TryEnterPump(aPump) then + try + aPump.ReportError(E); + finally + FCore.LeavePump; + end; + end; + end; + finally + FCore.WorkerDone; + end; + finally + try + FCore.ReleaseReference; + finally + CurrentWSPumpCore:=PreviousCore; + end; + end; + end; +end; + +{ TWebSocketClientConnection } + +procedure TWebSocketClientConnection.DoDisconnect; +begin + if Assigned(FClientSession) then + TWSClientSession(FClientSession).ConnectionRequestedDisconnect(Self); +end; + +procedure TWebSocketClientConnection.SetClientSession(aSession: TObject); +begin + FClientSession:=aSession; +end; + +procedure TWebSocketClientConnection.Send(aFrame: TWSFrame); +begin + if Assigned(FClientSession) and + (not TWSClientSession(FClientSession).IsOpen) then + Raise EWSReadInterrupted.Create(SErrClientSessionClosing); + if not HandshakeCompleted then + Raise EWebSocketClient.Create(SErrHandshakeInComplete); + inherited Send(aFrame); +end; + +function TWebSocketClientConnection.GetClient: TCustomWebsocketClient; + +begin + Result:=Owner as TCustomWebsocketClient; +end; + + +{ TCustomWebsocketClient } + +constructor TCustomWebsocketClient.Create(aOwner: TComponent); +begin + inherited Create(aOwner); + InitCriticalSection(FStateLock); + FMaxFramePayloadSize:=DefaultMaxFramePayloadSize; + FMaxMessagePayloadSize:=DefaultMaxMessagePayloadSize; + FOwnerGate:=TWSClientOwnerGate.Create(Self); +end; + +function TCustomWebsocketClient.GetActive: Boolean; +begin + EnterCriticalSection(FStateLock); + try + Result:=FActive; + finally + LeaveCriticalSection(FStateLock); + end; +end; + +function TCustomWebsocketClient.GetConnection: TWebSocketClientConnection; +begin + EnterCriticalSection(FStateLock); + try + Result:=FConnection; + finally + LeaveCriticalSection(FStateLock); + end; +end; + +function TCustomWebsocketClient.AcquireCurrentSession: TObject; +begin + EnterCriticalSection(FStateLock); + try + Result:=FSession; + if Assigned(Result) then + TWSClientSession(Result).AddReference; + finally + LeaveCriticalSection(FStateLock); + end; +end; + +function TCustomWebsocketClient.AcquireHandshakeSession: TObject; +begin + if CurrentWSHandshakeClient=Self then + Result:=CurrentWSHandshakeSession + else + Result:=Nil; + if Assigned(Result) then + TWSClientSession(Result).AddReference + else + Result:=AcquireCurrentSession; +end; + +function TCustomWebsocketClient.IsCurrentSession(aSession: TObject): Boolean; +begin + EnterCriticalSection(FStateLock); + try + Result:=(FSession=aSession) and (not FDestroying); + finally + LeaveCriticalSection(FStateLock); + end; +end; + +procedure TCustomWebsocketClient.DisconnectSession(aSession: TObject; + SendClose: Boolean; aEventSender: TObject); +var + Session : TWSClientSession; + Pump : TWSMessagePump; +begin + Session:=TWSClientSession(aSession); + if not Assigned(Session) then + Exit; + EnterCriticalSection(FStateLock); + try + Pump:=FMessagePump; + finally + LeaveCriticalSection(FStateLock); + end; + if SendClose and Session.IsOpen and Session.Connection.HandshakeCompleted and + (Session.Connection.CloseState<>csClosed) then + try + Session.Connection.Close(''); + except + on E : Exception do + if Assigned(Pump) then + Pump.ReportError(E); + end; + Session.NotifyDisconnected(Pump,aEventSender); +end; + +function TCustomWebsocketClient.DetachCurrentSession(aSession: TObject): Boolean; +begin + EnterCriticalSection(FStateLock); + try + Result:=Assigned(FSession) and + ((aSession=Nil) or (FSession=aSession)); + if Result then + begin + FSession:=Nil; + FConnection:=Nil; + FTransport:=Nil; + FSocket:=Nil; + FActive:=False; + end; + finally + LeaveCriticalSection(FStateLock); + end; +end; + +procedure TCustomWebsocketClient.SessionMessageReceived(aSession: TObject; + const aMessage: TWSMessage); +var + aHandler : TWSMessageEvent; +begin + aHandler:=FOnMessageReceived; + if Assigned(aHandler) then + aHandler(Self,aMessage); + { Do not access Self after user code. Disconnect/reconnect are supported; + direct Free from the callback is rejected by Destroy because a Pascal + destructor cannot safely return into this method. } +end; + +procedure TCustomWebsocketClient.SessionControlReceived(aSession: TObject; + aEventSender: TObject; aType: TFrameType; const aData: TBytes); +var + aHandler : TWSControlEvent; +begin + aHandler:=FOnControl; + if Assigned(aHandler) then + aHandler(aEventSender,aType,aData); + { As above, callback-time disconnect/reconnect are supported, not a direct + destruction of the callback target. } +end; + +procedure TCustomWebsocketClient.ReportCallbackError(aPump: TWSMessagePump; + E: Exception); +begin + if Assigned(aPump) then + aPump.ReportError(E); +end; + +procedure TCustomWebsocketClient.SessionDisconnected(aSession: TObject; + aEventSender: TObject; aPump: TWSMessagePump); +var + Session : TWSClientSession; + Pump : TWSMessagePump; + Handler : TNotifyEvent; + ReleaseClientReference : Boolean; +begin + Session:=TWSClientSession(aSession); + Pump:=aPump; + Session.BeginClosing; + EnterCriticalSection(FStateLock); + try + if FConnectingGeneration=Session.Generation then + FConnectingGeneration:=0; + ReleaseClientReference:=FSession=Session; + if ReleaseClientReference then + begin + FSession:=Nil; + FConnection:=Nil; + FTransport:=Nil; + FSocket:=Nil; + FActive:=False; + if not Assigned(Pump) then + Pump:=FMessagePump; + end; + Handler:=FOnDisconnect; + finally + LeaveCriticalSection(FStateLock); + end; + + { Logical removal precedes terminal cancellation and the application + notification. Neither registry nor component-state locks cross user + code. } + if Assigned(Pump) then + Pump.RemoveClient(Session.Connection); + Session.CancelReads; + if Assigned(Handler) then + try + Handler(aEventSender); + except + on E : Exception do + if Assigned(Pump) then + Pump.ReportError(E); + end; + if ReleaseClientReference then + try + Session.ReleaseReference; + except + on E : Exception do + if Assigned(Pump) then + Pump.ReportError(E); + end; + { Do not access Self after the handler. } +end; + +procedure TCustomWebsocketClient.SessionConnected(aSession: TObject); +var + Handler : TNotifyEvent; + IsCurrent : Boolean; +begin + EnterCriticalSection(FStateLock); + try + IsCurrent:=(FSession=aSession) and FActive and (not FDestroying); + Handler:=FOnConnect; + finally + LeaveCriticalSection(FStateLock); + end; + if IsCurrent and Assigned(Handler) then + Handler(Self); + { Do not access Self after the handler. } +end; + +procedure TCustomWebsocketClient.CheckInactive; +var + IsBusy : Boolean; +begin + EnterCriticalSection(FStateLock); + try + IsBusy:=FActive or Assigned(FSession) or + (FConnectingGeneration<>0) or FDestroying; + finally + LeaveCriticalSection(FStateLock); + end; + If IsBusy then + Raise EWebSocketClient.Create(SErrConnectionActive); +end; + +Function TCustomWebsocketClient.CheckIncoming : TIncomingResult; +var + Session : TWSClientSession; + Pump : TWSMessagePump; +begin + Pump:=MessagePump; + Session:=TWSClientSession(AcquireCurrentSession); + if not Assigned(Session) then + Raise EWebSocketClient.Create(SErrConnectionInActive); + try + if not Session.Connection.HandshakeCompleted then + Raise EWebSocketClient.Create(SErrHandshakeInComplete); + try + Result:=Session.Connection.CheckIncoming(CheckTimeout); + if Result=irClose then + Session.NotifyDisconnected(Pump,Session.Connection); + except + { Manual polling gets the same terminal lifetime transition as the + threaded pump. The original exception still reaches the caller. } + on E : Exception do + begin + Session.NotifyDisconnected(Pump,Session.Connection); + raise; + end; + end; + finally + ReleaseClientSession(Session,Pump); + end; +end; + +procedure TCustomWebsocketClient.ControlReceived(Sender: TObject; aType : TFrameType; const aData: TBytes); +begin + If Assigned(FOnControl) then + FOnControl(Sender, aType, aData); +end; + +function TCustomWebsocketClient.CreateClientConnection(aTransport: TWSClientTRansport): TWebsocketClientConnection; + +begin + Result:=TWebSocketClientConnection.Create(Self,aTransport,FOptions); +end; + +procedure TCustomWebsocketClient.ConnectionDisconnected(Sender : TObject); +var + Session : TWSClientSession; + Pump : TWSMessagePump; +begin + Pump:=MessagePump; + Session:=TWSClientSession(AcquireCurrentSession); + if not Assigned(Session) then + Exit; + try + Session.ConnectionRequestedDisconnect(Sender); + finally + ReleaseClientSession(Session,Pump); + end; +end; + +procedure TCustomWebsocketClient.Connect; +var + SSLHandler: TSSLSocketHandler; + NewSocket : TInetSocket; + NewTransport : TWSClientTransport; + NewConnection : TWebSocketClientConnection; + Session : TWSClientSession; + Gate : TWSClientOwnerGate; + Generation : QWord; + Pump : TWSMessagePump; + Published : Boolean; + HandshakeOK : Boolean; + FinishConnection : Boolean; + CallbackClient : TCustomWebsocketClient; + ConnectToken : TWSOwnerGateToken; + PreviousHandshakeClient : TCustomWebsocketClient; + PreviousHandshakeSession : TWSClientSession; +begin + EnterCriticalSection(FStateLock); + try + if FDestroying then + Raise EWebSocketClient.Create(SErrConnectionInActive); + if Assigned(FSession) or FActive or (FConnectingGeneration<>0) then + Exit; + Inc(FNextGeneration); + if FNextGeneration=0 then + Inc(FNextGeneration); + Generation:=FNextGeneration; + FConnectingGeneration:=Generation; + Gate:=TWSClientOwnerGate(FOwnerGate); + Gate.AddReference; // protects the handoff into the first session lease + if not Gate.TryEnter(CallbackClient,ConnectToken) then + begin + FConnectingGeneration:=0; + Gate.ReleaseReference; + Raise EWebSocketClient.Create(SErrConnectionInActive); + end; + finally + LeaveCriticalSection(FStateLock); + end; + + NewSocket:=Nil; + NewTransport:=Nil; + NewConnection:=Nil; + Session:=Nil; + Published:=False; + Pump:=Nil; + SSLHandler := nil; + try + Session:=TWSClientSession.Create(Gate,Generation); + if UseSSL then + begin + SSLHandler := TSSLSocketHandler.GetDefaultHandler; + SSLHandler.VerifyPeerCert := False; + end; + NewSocket:=TInetSocket.Create(HostName,Port,ConnectTimeout,SSLHandler); + NewTransport:=TWSClientTransport.Create(NewSocket); + NewSocket:=Nil; // owned by NewTransport + NewConnection:=CreateClientConnection(NewTransport); + if not Assigned(NewConnection) then + Raise EWebSocketClient.Create(SErrConnectionInActive); + NewTransport:=Nil; // owned by NewConnection + Session.AttachConnection(NewConnection); + NewConnection:=Nil; // owned by Session + + if OutGoingFrameMask=0 then + OutGoingFrameMask:=1+Random(MaxInt-1); + Session.Connection.OutgoingFrameMask:=OutGoingFrameMask; + Session.Connection.MaxFramePayloadSize:=FMaxFramePayloadSize; + Session.Connection.MaxMessagePayloadSize:=FMaxMessagePayloadSize; + if UseSSL then + TInetSocket(Session.Connection.ClientTransport.Socket).Connect; + + EnterCriticalSection(FStateLock); + try + if FDestroying or Assigned(FSession) or + (FConnectingGeneration<>Generation) then + Raise EWebSocketClient.Create(SErrConnectionActive); + Session.AddReference; // current-client ownership + FSession:=Session; + FConnection:=Session.Connection; + FTransport:=Session.Connection.ClientTransport; + FSocket:=FTransport.Socket as TInetSocket; + FActive:=True; + Pump:=FMessagePump; + Published:=True; + finally + LeaveCriticalSection(FStateLock); + end; + + try + PreviousHandshakeClient:=CurrentWSHandshakeClient; + PreviousHandshakeSession:=CurrentWSHandshakeSession; + CurrentWSHandshakeClient:=Self; + CurrentWSHandshakeSession:=Session; + try + HandshakeOK:=DoHandShake; + finally + CurrentWSHandshakeSession:=PreviousHandshakeSession; + CurrentWSHandshakeClient:=PreviousHandshakeClient; + end; + if (not HandshakeOK) or (not IsCurrentSession(Session)) or + (not Session.IsOpen) then + begin + DisconnectSession(Session,False,Self); + Exit; + end; + + EnterCriticalSection(FStateLock); + try + FinishConnection:=(FSession=Session) and Session.IsOpen and + (not FDestroying); + if FinishConnection then + Pump:=FMessagePump; + finally + LeaveCriticalSection(FStateLock); + end; + if not FinishConnection then + begin + DisconnectSession(Session,False,Self); + Exit; + end; + + if Assigned(Pump) then + Pump.AddClient(Session.Connection); + + FinishConnection:=IsCurrentSession(Session) and Session.IsOpen; + EnterCriticalSection(FStateLock); + try + FinishConnection:=FinishConnection and (FSession=Session) and + (not FDestroying); + if FinishConnection and (FConnectingGeneration=Generation) then + FConnectingGeneration:=0; + finally + LeaveCriticalSection(FStateLock); + end; + if not FinishConnection then + begin + if Assigned(Pump) then + Pump.RemoveClient(Session.Connection); + DisconnectSession(Session,False,Self); + Exit; + end; + Session.NotifyConnected; + except + DisconnectSession(Session,False,Self); + raise; + end; + finally + try + EnterCriticalSection(FStateLock); + try + if FConnectingGeneration=Generation then + FConnectingGeneration:=0; + finally + LeaveCriticalSection(FStateLock); + end; + try + if not Published then + begin + try + NewConnection.Free; + finally + try + NewTransport.Free; + finally + NewSocket.Free; + end; + end; + end; + finally + ReleaseClientSession(Session,Pump); + end; + finally + try + Gate.Leave(ConnectToken); + finally + Gate.ReleaseReference; + end; + end; + end; +end; + + +destructor TCustomWebsocketClient.Destroy; +var + Session : TWSClientSession; + Gate : TWSClientOwnerGate; + Pump : TWSMessagePump; +begin + Session:=Nil; + Gate:=TWSClientOwnerGate(FOwnerGate); + if Gate.IsCurrentCallback then + Raise EWebSocketClient.Create(SErrFreeFromCallback); + EnterCriticalSection(FStateLock); + try + FDestroying:=True; + if Assigned(FSession) then + begin + Session:=TWSClientSession(FSession); + FSession:=Nil; + end; + FConnection:=Nil; + FTransport:=Nil; + FSocket:=Nil; + FActive:=False; + Pump:=FMessagePump; + finally + LeaveCriticalSection(FStateLock); + end; + try + if Assigned(Session) then + begin + Session.BeginClosing; + if Assigned(Pump) then + Pump.RemoveClient(Session.Connection); + Session.CancelReads; + Session.NotifyDisconnected(Pump,Self); + end; + Gate.CloseAdmission; + Gate.WaitForQuiescence; + if Assigned(Session) then + try + Session.ReleaseReference; + except + on E : Exception do + if Assigned(Pump) then + Pump.ReportError(E); + end; + FreeAndNil(FHandShake); + FreeAndNil(FHandshakeResponse); + if Assigned(FMessagePump) then + FMessagePump.RemoveFreeNotification(Self); + FMessagePump:=Nil; + Gate.ReleaseReference; + FOwnerGate:=Nil; + finally + DoneCriticalSection(FStateLock); + inherited Destroy; + end; +end; + + +Function TCustomWebsocketClient.CreateHandShakeRequest : TWSHandShakeRequest; + +begin + Result:=TWSHandShakeRequest.Create('',Nil); +end; + +procedure TCustomWebsocketClient.SendData(aBytes: TBytes); +var + Session : TWSClientSession; + Pump : TWSMessagePump; +begin + Pump:=MessagePump; + Session:=TWSClientSession(AcquireCurrentSession); + if not Assigned(Session) then + Raise EWebSocketClient.Create(SErrConnectionInActive); + try + Session.Connection.Send(aBytes); + finally + ReleaseClientSession(Session,Pump); + end; +end; + +procedure TCustomWebsocketClient.SendHeaders(aHeaders : TStrings); + +Var + HeaderBlock : String; + B : TBytes; + Session : TWSClientSession; + Pump : TWSMessagePump; + I : Integer; + +begin + Pump:=MessagePump; + Session:=TWSClientSession(AcquireHandshakeSession); + if not Assigned(Session) then + Raise EWebSocketClient.Create(SErrConnectionInActive); + try + if (not Session.IsOpen) or (not IsCurrentSession(Session)) then + Exit; + HeaderBlock:=''; + for I:=0 to aHeaders.Count-1 do + HeaderBlock:=HeaderBlock+aHeaders[I]+#13#10; + HeaderBlock:=HeaderBlock+#13#10; + B:=TEncoding.UTF8.GetAnsiBytes(HeaderBlock); + { One write-all call keeps a websocket frame from being interleaved into + the HTTP upgrade request. } + Session.Connection.Transport.WriteBuffer(B); + finally + ReleaseClientSession(Session,Pump); + end; +end; + +procedure TCustomWebsocketClient.SendHandShakeRequest; + +Var + aRequest : TWSHandShakeRequest; + aHeaders : TStrings; + Session : TWSClientSession; + Pump : TWSMessagePump; + CallbackClient : TCustomWebsocketClient; + Handler : TWSClientHandshakeEvent; + GateToken : TWSOwnerGateToken; +begin + Pump:=MessagePump; + Session:=TWSClientSession(AcquireHandshakeSession); + if not Assigned(Session) then + Raise EWebSocketClient.Create(SErrConnectionInActive); + aHeaders:=Nil; + aRequest:=Nil; + try + if not Session.FOwnerGate.TryEnter(CallbackClient,GateToken) then + Exit; + try + if (not Session.IsOpen) or + (not CallbackClient.IsCurrentSession(Session)) then + Exit; + aRequest:=CallbackClient.CreateHandShakeRequest; + if not Assigned(aRequest) then + Raise EWebSocketClient.Create(SErrHandshakeInComplete); + if (not Session.IsOpen) or + (not CallbackClient.IsCurrentSession(Session)) then + Exit; + aRequest.Host:=HostName; + aRequest.Port:=Port; + aRequest.Resource:=Resource; + aHeaders:=TStringList.Create; + aHeaders.NameValueSeparator:=':'; + aRequest.ToStrings(aHeaders); + Session.SetHandshakeRequest(aRequest); + aRequest:=Nil; + Handler:=CallbackClient.FOnSendHandshake; + if Assigned(Handler) then + Handler(CallbackClient,aHeaders); + if (not Session.IsOpen) or + (not CallbackClient.IsCurrentSession(Session)) then + Exit; + // Do not use FClient.WriteHeader, it messes up the strings ! + CallbackClient.SendHeaders(aHeaders); + finally + Session.FOwnerGate.Leave(GateToken); + end; + finally + aHeaders.Free; + aRequest.Free; + ReleaseClientSession(Session,Pump); + end; +end; + +procedure TCustomWebsocketClient.SendMessage(const aMessage: String); +var + Session : TWSClientSession; + Pump : TWSMessagePump; +begin + Pump:=MessagePump; + Session:=TWSClientSession(AcquireCurrentSession); + if not Assigned(Session) then + Raise EWebSocketClient.Create(SErrConnectionInActive); + try + Session.Connection.Send(aMessage); + finally + ReleaseClientSession(Session,Pump); + end; +end; + +Function TCustomWebsocketClient.CreateHandshakeResponse(aHeaders : TStrings) : TWSHandShakeResponse; + +begin + Result:=TWSHandShakeResponse.Create('',aHeaders); +end; + +Function TCustomWebsocketClient.CheckHandShakeResponse(aHeaders : TStrings) : Boolean; + + Function ParseStatusLine(const aLine : String; + out aHTTPVersion : String; out aStatusCode : Integer; + out aStatusText : String) : Boolean; + Var + P : Integer; + ProtocolPart, + Rest, + CodePart : String; + begin + Result:=False; + aHTTPVersion:=''; + aStatusCode:=0; + aStatusText:=''; + Rest:=Trim(aLine); + P:=Pos(' ',Rest); + if P=0 then + Exit; + ProtocolPart:=Copy(Rest,1,P-1); + if (Length(ProtocolPart)<5) or + (CompareText(Copy(ProtocolPart,1,5),'HTTP/')<>0) then + Exit; + aHTTPVersion:=Copy(ProtocolPart,6,MaxInt); + if aHTTPVersion='' then + Exit; + Rest:=Trim(Copy(Rest,P+1,MaxInt)); + P:=Pos(' ',Rest); + if P=0 then + begin + CodePart:=Rest; + aStatusText:=''; + end + else + begin + CodePart:=Copy(Rest,1,P-1); + aStatusText:=Trim(Copy(Rest,P+1,MaxInt)); + end; + Result:=TryStrToInt(CodePart,aStatusCode); + end; + + Function HasHeaderToken(const aValue,aToken : String) : Boolean; + Var + P : Integer; + Remaining, + ValuePart : String; + begin + Remaining:=aValue; + repeat + P:=Pos(',',Remaining); + if P=0 then + begin + ValuePart:=Trim(Remaining); + Remaining:=''; + end + else + begin + ValuePart:=Trim(Copy(Remaining,1,P-1)); + Delete(Remaining,1,P); + end; + if SameText(ValuePart,aToken) then + Exit(True); + until Remaining=''; + Result:=False; + end; + +Var + K : String; + {%H-}hash : TSHA1Digest; + B : TBytes; + Session : TWSClientSession; + Pump : TWSMessagePump; + Response : TWSHandShakeResponse; + ValidStatus : Boolean; + HTTPVersion, + StatusText : String; + StatusCode : Integer; + +begin + Pump:=MessagePump; + Session:=TWSClientSession(AcquireHandshakeSession); + if not Assigned(Session) then + Raise EWebSocketClient.Create(SErrConnectionInActive); + B:=[]; + Response:=Nil; + try + if not Assigned(Session.HandshakeRequest) then + Raise EWebSocketClient.Create(SErrHandshakeInComplete); + ValidStatus:=False; + if aHeaders.Count>0 then + ValidStatus:=ParseStatusLine(aHeaders[0],HTTPVersion,StatusCode, + StatusText) and (StatusCode=101); + Response:=CreateHandshakeResponse(aHeaders); + if not Assigned(Response) then + Raise EWebSocketClient.Create(SErrHandshakeInComplete); + Response.HTTPVersion:=HTTPVersion; + Response.StatusCode:=StatusCode; + Response.StatusText:=StatusText; + k := Trim(Session.HandshakeRequest.Key) + SSecWebSocketGUID; + hash:=SHA1String(k); + SetLength(B,SizeOf(hash)); + Move(hash[0],B[0],SizeOf(hash)); + k:=EncodeBytesBase64(B); + { Sec-WebSocket-Accept is base64 and therefore case-sensitive. } + Result:=(K=Response.Accept) + and SameText(Response.Upgrade,'websocket') + and HasHeaderToken(Response.Connection,'Upgrade') + and ValidStatus; + Session.SetHandshakeResponse(Response); + Response:=Nil; + finally + Response.Free; + ReleaseClientSession(Session,Pump); + end; +end; + +Function TCustomWebsocketClient.ReadHandShakeResponse : Boolean; + + Function ParseResponseStatus(const aLine : String; + out aHTTPVersion : String; out aStatusCode : Integer; + out aStatusText : String) : Boolean; + Var + P : Integer; + ProtocolPart, + Rest, + CodePart : String; + begin + Result:=False; + aHTTPVersion:=''; + aStatusCode:=0; + aStatusText:=''; + Rest:=Trim(aLine); + P:=Pos(' ',Rest); + if P=0 then + Exit; + ProtocolPart:=Copy(Rest,1,P-1); + if (Length(ProtocolPart)<6) or + (CompareText(Copy(ProtocolPart,1,5),'HTTP/')<>0) then + Exit; + aHTTPVersion:=Copy(ProtocolPart,6,MaxInt); + Rest:=Trim(Copy(Rest,P+1,MaxInt)); + P:=Pos(' ',Rest); + if P=0 then + CodePart:=Rest + else + begin + CodePart:=Copy(Rest,1,P-1); + aStatusText:=Trim(Copy(Rest,P+1,MaxInt)); + end; + Result:=(aHTTPVersion<>'') and TryStrToInt(CodePart,aStatusCode); + end; + +Var + S : String; + aHeaders : TStrings; + aResponse : TWSHandShakeResponse; + ResponseToTransfer : TWSHandShakeResponse; + aHandler : TWSClientHandshakeResponseEvent; + Session : TWSClientSession; + Pump : TWSMessagePump; + CallbackClient : TCustomWebsocketClient; + GateToken : TWSOwnerGateToken; + HeaderBytes : SizeInt; + HeaderLines : Integer; + LineBytes : SizeInt; + HTTPVersion, + StatusText : String; + StatusCode : Integer; + ValidStatus : Boolean; + +begin + Result:=False; + Pump:=MessagePump; + Session:=TWSClientSession(AcquireHandshakeSession); + if not Assigned(Session) then + Raise EWebSocketClient.Create(SErrConnectionInActive); + aHeaders:=TStringList.Create; + ResponseToTransfer:=Nil; + HeaderBytes:=0; + HeaderLines:=0; + HTTPVersion:=''; + StatusText:=''; + StatusCode:=0; + Try + if not Session.FOwnerGate.TryEnter(CallbackClient,GateToken) then + Exit; + try + if (not Session.IsOpen) or + (not CallbackClient.IsCurrentSession(Session)) then + Exit; + aHeaders.NameValueSeparator:=':'; + Repeat + S:=Session.Connection.Transport.ReadLn; + LineBytes:=Length(S); + if LineBytes>WSMaxHandshakeHeaderLineBytes then + Raise EWSHandShake.Create(SErrHandshakeHeaderLineTooLong); + Inc(HeaderLines); + if (HeaderLines>WSMaxHandshakeHeaderLines) or + (LineBytes>WSMaxHandshakeHeaderBytes-HeaderBytes-2) then + Raise EWSHandShake.Create(SErrHandshakeHeadersTooLarge); + Inc(HeaderBytes,LineBytes+2); + aHeaders.Add(S); + Until (S=''); + ValidStatus:=(aHeaders.Count>0) and + ParseResponseStatus(aHeaders[0],HTTPVersion,StatusCode,StatusText) and + (StatusCode=101); + Result:=ValidStatus and CallbackClient.CheckHandShakeResponse(aHeaders); + if (not Session.IsOpen) or + (not CallbackClient.IsCurrentSession(Session)) then + begin + Result:=False; + Exit; + end; + if Result then + begin + ResponseToTransfer:=Session.TakeHandshakeResponse; + if not Assigned(ResponseToTransfer) then + ResponseToTransfer:=CallbackClient.CreateHandshakeResponse(aHeaders); + if (not Assigned(ResponseToTransfer)) or (not Session.IsOpen) or + (not CallbackClient.IsCurrentSession(Session)) then + begin + Result:=False; + Exit; + end; + { Even an override which implements its own header checks exposes the + status line actually received, rather than constructor defaults. } + ResponseToTransfer.HTTPVersion:=HTTPVersion; + ResponseToTransfer.StatusCode:=StatusCode; + ResponseToTransfer.StatusText:=StatusText; + { Keep the original public non-owning HandshakeResponse contract. + The generation session owns this response for exactly as long as the + connection may expose it. } + Session.SetHandshakeResponse(ResponseToTransfer); + ResponseToTransfer:=Nil; + aResponse:=Session.FHandshakeResponse; + Session.Connection.HandshakeResponse:=aResponse; + aHandler:=CallbackClient.FOnHandshakeResponse; + if Assigned(aHandler) then + aHandler(CallbackClient,aResponse,Result); + end; + finally + Session.FOwnerGate.Leave(GateToken); + end; + Finally + ResponseToTransfer.Free; + aHeaders.Free; + ReleaseClientSession(Session,Pump); + End; +end; + +Function TCustomWebsocketClient.DoHandShake : Boolean; +var + Session : TWSClientSession; + PreviousSession : TWSClientSession; + PreviousClient : TCustomWebsocketClient; + Pump : TWSMessagePump; + CallbackClient : TCustomWebsocketClient; + GateToken : TWSOwnerGateToken; +begin + Result:=False; + Pump:=MessagePump; + Session:=TWSClientSession(AcquireHandshakeSession); + if not Assigned(Session) then + Raise EWebSocketClient.Create(SErrConnectionInActive); + try + if not Session.FOwnerGate.TryEnter(CallbackClient,GateToken) then + Exit; + try + PreviousSession:=CurrentWSHandshakeSession; + PreviousClient:=CurrentWSHandshakeClient; + CurrentWSHandshakeClient:=CallbackClient; + CurrentWSHandshakeSession:=Session; + try + if (not Session.IsOpen) or + (not CallbackClient.IsCurrentSession(Session)) then + Exit; + CallbackClient.SendHandShakeRequest; + if (not Session.IsOpen) or + (not CallbackClient.IsCurrentSession(Session)) then + Exit; + Result:=CallbackClient.ReadHandShakeResponse; + Result:=Result and Session.IsOpen and + CallbackClient.IsCurrentSession(Session); + finally + CurrentWSHandshakeSession:=PreviousSession; + CurrentWSHandshakeClient:=PreviousClient; + end; + finally + Session.FOwnerGate.Leave(GateToken); + end; + finally + ReleaseClientSession(Session,Pump); + end; +end; + +procedure TCustomWebsocketClient.Loaded; +begin + inherited; + if FLoadActive then + Connect; +end; + +procedure TCustomWebsocketClient.Notification(aComponent : TComponent; + Operation : TOperation); +begin + inherited Notification(aComponent,Operation); + if Operation=opRemove then + begin + EnterCriticalSection(FStateLock); + try + if aComponent=FMessagePump then + FMessagePump:=Nil; + finally + LeaveCriticalSection(FStateLock); + end; + end; +end; + +procedure TCustomWebsocketClient.MessageReceived(Sender: TObject; const aMessage : TWSMessage) ; +begin + if Assigned(OnMessageReceived) and (TWSClientConnection(Sender).HandshakeCompleted) then + OnMessageReceived(Self, AMessage); +end; + +procedure TCustomWebsocketClient.Ping(aMessage: UTF8String); +var + Session : TWSClientSession; + Pump : TWSMessagePump; +begin + Pump:=MessagePump; + Session:=TWSClientSession(AcquireCurrentSession); + if not Assigned(Session) then + Raise EWebSocketClient.Create(SErrConnectionInActive); + try + Session.Connection.Send(ftPing,TEncoding.UTF8.GetAnsiBytes(aMessage)); + finally + ReleaseClientSession(Session,Pump); + end; +end; + +procedure TCustomWebsocketClient.Pong(aMessage: UTF8String); +var + Session : TWSClientSession; + Pump : TWSMessagePump; +begin + Pump:=MessagePump; + Session:=TWSClientSession(AcquireCurrentSession); + if not Assigned(Session) then + Raise EWebSocketClient.Create(SErrConnectionInActive); + try + Session.Connection.Send(ftPong,TEncoding.UTF8.GetAnsiBytes(aMessage)); + finally + ReleaseClientSession(Session,Pump); + end; +end; + +procedure TCustomWebsocketClient.FreeConnectionObjects; +var + Session : TWSClientSession; + Pump : TWSMessagePump; +begin + Pump:=MessagePump; + Session:=TWSClientSession(AcquireCurrentSession); + if not Assigned(Session) then + Exit; + try + Disconnect(False); + finally + ReleaseClientSession(Session,Pump); + end; +end; + +procedure TCustomWebsocketClient.Disconnect(SendClose : boolean = true); +var + Session : TWSClientSession; + Pump : TWSMessagePump; +begin + Pump:=MessagePump; + Session:=TWSClientSession(AcquireCurrentSession); + if not Assigned(Session) then + Exit; + try + { SessionDisconnected transfers/releases the client's ownership. This + local reference keeps the session alive through the whole operation. } + DisconnectSession(Session,SendClose,Self); + finally + ReleaseClientSession(Session,Pump); + end; +end; + +procedure TCustomWebsocketClient.SetActive(const Value: Boolean); +begin + FLoadActive := Value; + if (csDesigning in ComponentState) then + exit; + if Value then + Connect + else + Disconnect; +end; + +procedure TCustomWebsocketClient.SetAutoCheckMessages(const Value: Boolean); +begin + CheckInactive; + FAutoCheckMessages := Value; +end; + +procedure TCustomWebsocketClient.SetCheckTimeOut(const Value: Integer); +begin + CheckInactive; + FCheckTimeOut := Value; +end; + +procedure TCustomWebsocketClient.SetConnectTimeout(const Value: Integer); +begin + CheckInactive; + FConnectTimeout := Value; +end; + +procedure TCustomWebsocketClient.SetHostName(const Value: String); +begin + CheckInactive; + FHostName := Value; +end; + +procedure TCustomWebsocketClient.SetMessagePump(AValue: TWSMessagePump); +begin + if FMessagePump=AValue then Exit; + if Active or TWSClientOwnerGate(FOwnerGate).HasCallbacks then + Raise EWebSocketClient.Create(SErrConnectionActive); + If Assigned(FMessagePump) then + FMessagePump.RemoveFreeNotification(Self); + FMessagePump:=AValue; + If Assigned(FMessagePump) then + FMessagePump.FreeNotification(Self); +end; + +procedure TCustomWebsocketClient.SetOptions(const Value: TWSOptions); +begin + CheckInactive; + FOptions := Value; +end; + +procedure TCustomWebsocketClient.SetMaxFramePayloadSize(const Value: QWord); +begin + CheckInactive; + FMaxFramePayloadSize:=Value; +end; + +procedure TCustomWebsocketClient.SetMaxMessagePayloadSize(const Value: QWord); +begin + CheckInactive; + FMaxMessagePayloadSize:=Value; +end; + +procedure TCustomWebsocketClient.SetPort(const Value: Integer); +begin + CheckInactive; + FPort := Value; +end; + +procedure TCustomWebsocketClient.SetResource(const Value: string); +begin + CheckInactive; + FResource := Value; +end; + +procedure TCustomWebsocketClient.SetUseSSL(const Value: Boolean); +begin + CheckInactive; + FUseSSL := Value; end; @@ -576,19 +2443,272 @@ procedure TCustomWebsocketClient.SetUseSSL(const Value: Boolean); { TWSMessagePump } procedure TWSMessagePump.AddClient(aConnection: TWSClientConnection); +var + Session : TWSClientSession; + Entry : TWSMessagePumpEntry; + Entries : TList; + I : Integer; begin - List.Add(aConnection); + { A raw TWSClientConnection has no reference/lifetime protocol. Accepting + one here would let its external owner free it while CheckIncoming runs. + Preserve the public entry point, but fail before registration unless the + connection carries the managed session used by this client unit. } + if not (aConnection is TWebSocketClientConnection) then + Raise EWebSocketClient.Create('The message pump requires a client session connection'); + Session:=TWSClientSession(TWebSocketClientConnection(aConnection).FClientSession); + if not Assigned(Session) then + Raise EWebSocketClient.Create('The websocket connection has no client session'); + + Entry:=Nil; + EnterCriticalSection(FRegistryLock); + try + Entries:=FEntries.LockList; + try + for I:=0 to Entries.Count-1 do + if TWSMessagePumpEntry(Entries[I]).Connection=aConnection then + Exit; + Entry:=TWSMessagePumpEntry.Create(Session,aConnection); + Entries.Add(Entry); + Entry.AddReference; // registry lease; initial reference stays local + finally + FEntries.UnlockList; + end; + try + FList.Add(aConnection); // protected compatibility mirror + except + Entries:=FEntries.LockList; + try + Entries.Remove(Entry); + Entry.MarkRemoved; + finally + FEntries.UnlockList; + end; + try + Entry.ReleaseReference; // registry lease + except + on Exception do + ; // preserve the original AddClient failure + end; + try + Entry.ReleaseReference; // local lease + except + on Exception do + ; // preserve the original AddClient failure + end; + Entry:=Nil; + raise; + end; + finally + LeaveCriticalSection(FRegistryLock); + end; + try + { The default thread pump schedules this entry through its virtual + CheckConnections/ReadConnections pipeline. That retains the protected + extension hook without introducing a second socket reader. } + if (not Entry.IsRegistered) or (not Session.IsOpen) then + RemoveEntry(Entry); + finally + try + Entry.ReleaseReference; // local lease + except + on E : Exception do + ReportError(E); + end; + end; end; procedure TWSMessagePump.RemoveClient(aConnection: TWSClientConnection); +var + Entry : TWSMessagePumpEntry; +begin + Entry:=TWSMessagePumpEntry(FindEntry(aConnection)); + if not Assigned(Entry) then + begin + FList.Remove(aConnection); + Exit; + end; + try + RemoveEntry(Entry); + finally + try + Entry.ReleaseReference; + except + on E : Exception do + ReportError(E); + end; + end; +end; + +function TWSMessagePump.FindEntry(aConnection: TWSClientConnection): TObject; +var + Entries : TList; + Entry : TWSMessagePumpEntry; + I : Integer; +begin + Result:=Nil; + EnterCriticalSection(FRegistryLock); + try + Entries:=FEntries.LockList; + try + for I:=0 to Entries.Count-1 do + begin + Entry:=TWSMessagePumpEntry(Entries[I]); + if Entry.Connection=aConnection then + begin + Entry.AddReference; + Result:=Entry; + Exit; + end; + end; + finally + FEntries.UnlockList; + end; + finally + LeaveCriticalSection(FRegistryLock); + end; +end; + +procedure TWSMessagePump.RemoveEntry(aEntry: TObject); +var + Entry : TWSMessagePumpEntry; + Entries : TList; + RemovedRegistryReference : Boolean; +begin + Entry:=TWSMessagePumpEntry(aEntry); + RemovedRegistryReference:=False; + EnterCriticalSection(FRegistryLock); + try + Entries:=FEntries.LockList; + try + if Entry.MarkRemoved then + begin + if Entries.Remove(Entry)>=0 then + RemovedRegistryReference:=True; + end; + finally + FEntries.UnlockList; + end; + FList.Remove(Entry.Connection); + finally + LeaveCriticalSection(FRegistryLock); + end; + { This is terminal cancellation for this exact session. It remains set + between exact-read chunks and therefore has no interrupt gap. } + Entry.CancelReads; + if RemovedRegistryReference then + try + Entry.ReleaseReference; + except + on E : Exception do + ReportError(E); + end; +end; + +procedure TWSMessagePump.StartEntry(aEntry: TObject; + aRunGeneration: LongInt); +var + Entry : TWSMessagePumpEntry; + Reader : TWSClientReaderThread; +begin + if not (Self is TWSThreadMessagePump) then + Exit; + if not TWSPumpCore(FCore).IsRunning(aRunGeneration) then + Exit; + Entry:=TWSMessagePumpEntry(aEntry); + if not Entry.TryStartWorker then + Exit; + Reader:=Nil; + try + Reader:=TWSClientReaderThread.Create(TWSPumpCore(FCore),Entry, + aRunGeneration); + Reader.Start; + except + if Assigned(Reader) then + begin + Reader.FreeOnTerminate:=False; + Reader.AbandonBeforeStart; + Reader.Free; + end + else + Entry.WorkerStopped; + raise; + end; +end; + +procedure TWSMessagePump.EntryEnded(aEntry: TObject; aError: Exception); +var + Entry : TWSMessagePumpEntry; +begin + Entry:=TWSMessagePumpEntry(aEntry); + Entry.Session.BeginClosing; + RemoveEntry(Entry); + if Assigned(aError) then + ReportError(aError); + { Notification is the tail operation: it may re-enter or attempt to free + its pump. } + Entry.Session.NotifyDisconnected(Self,Entry.Connection); +end; + +procedure TWSMessagePump.ReportError(aError: Exception); +var + Handler : TWSErrorEvent; +begin + Handler:=FOnError; + if Assigned(Handler) then + try + Handler(Self,aError); + except + { An error reporter must not terminate a reader or skip its releases. } + end; +end; + +procedure TWSMessagePump.ClearEntries; +var + Entries : TList; + LocalEntries : TList; + Entry : TWSMessagePumpEntry; + I : Integer; begin - FList.Remove(aConnection); + LocalEntries:=TList.Create; + try + EnterCriticalSection(FRegistryLock); + try + Entries:=FEntries.LockList; + try + for I:=0 to Entries.Count-1 do + begin + Entry:=TWSMessagePumpEntry(Entries[I]); + Entry.MarkRemoved; + LocalEntries.Add(Entry); + end; + Entries.Clear; + finally + FEntries.UnlockList; + end; + FList.Clear; + finally + LeaveCriticalSection(FRegistryLock); + end; + for I:=0 to LocalEntries.Count-1 do + begin + Entry:=TWSMessagePumpEntry(LocalEntries[I]); + try + Entry.ReleaseReference; + except + on E : Exception do + ReportError(E); + end; + end; + finally + LocalEntries.Free; + end; end; procedure TWSMessagePump.SetInterval(AValue: Integer); begin if FInterval=AValue then Exit; FInterval:=AValue; + TWSPumpCore(FCore).SetInterval(aValue); end; Function TWSMessagePump.WaitForData : Boolean; @@ -618,151 +2738,380 @@ procedure TWSMessagePump.SetInterval(AValue: Integer); end; function TWSMessagePump.CheckConnections: Boolean; - Var - aList : TList; - aClient: TWSClientConnection; - aTrans : TWSClientTransport; - I,aLen : Integer; + Entries : TList; begin - Result:=False; - aList := List.LockList; + Entries:=FEntries.LockList; try - aLen:=0; - SetLength(FReads,aList.Count); - for I := 0 to aList.Count - 1 do - begin - aClient := TWSClientConnection(aList.Items[I]); - if assigned(aClient) then - aTrans:=aClient.ClientTransport - else - aTrans:=Nil; - if (aTrans<>nil) then - begin - // There is already data - FReads[aLen]:=aTrans.Socket; - Inc(aLen); - end; - end; + Result:=Entries.Count<>0; finally - List.UnlockList; + FEntries.UnlockList; end; - if Not Result then - Result:=WaitForData; + if not Result then + TThread.Sleep(FInterval); end; constructor TWSMessagePump.Create(aOwner : TComponent); begin + inherited Create(aOwner); + InitCriticalSection(FRegistryLock); + FEntries:=TThreadList.Create; FList:=TThreadList.Create; FReads:=[]; FExceptions:=[]; Finterval:=25; + FCore:=TWSPumpCore.Create(Self,FInterval); end; destructor TWSMessagePump.Destroy; begin + TWSPumpCore(FCore).RequestStop; + TWSPumpCore(FCore).WaitWorkers(-1); + ClearEntries; + TWSPumpCore(FCore).CloseOwner; + TWSPumpCore(FCore).WaitOwnerUsers; + TWSPumpCore(FCore).ReleaseReference; + FCore:=Nil; + FreeAndNil(FEntries); FreeAndNil(FList); + DoneCriticalSection(FRegistryLock); inherited; end; -procedure TWSMessagePump.ReadConnections; +procedure TWSMessagePump.InterruptConnections; +Var + Entries : TList; + Snapshot : TList; + Entry : TWSMessagePumpEntry; + I : Integer; + +begin + Snapshot:=TList.Create; + try + Entries:=FEntries.LockList; + try + for I:=0 to Entries.Count-1 do + begin + Entry:=TWSMessagePumpEntry(Entries[I]); + Entry.AddReference; + Snapshot.Add(Entry); + end; + finally + FEntries.UnlockList; + end; + for I:=0 to Snapshot.Count-1 do + begin + Entry:=TWSMessagePumpEntry(Snapshot[I]); + try + try + if Entry.WorkerRunning then + Entry.InterruptRead; + except + on E : Exception do + ReportError(E); + end; + finally + try + Entry.ReleaseReference; + except + on E : Exception do + ReportError(E); + end; + end; + end; + finally + Snapshot.Free; + end; +end; +procedure TWSMessagePump.ReadConnections; Var - aList : TList; - aClient: TWSClientConnection; + Entries : TList; + Snapshot : TList; + Entry : TWSMessagePumpEntry; + IncomingResult: TIncomingResult; + RunGeneration : LongInt; I : Integer; begin + if Self is TWSThreadMessagePump then + begin + { Persistent per-session readers are the sole consumers in the default + threaded pump. The legacy protected pipeline schedules them here. } + Snapshot:=TList.Create; + try + Entries:=FEntries.LockList; + try + for I:=0 to Entries.Count-1 do + begin + Entry:=TWSMessagePumpEntry(Entries[I]); + Entry.AddReference; + Snapshot.Add(Entry); + end; + finally + FEntries.UnlockList; + end; + RunGeneration:=TWSPumpCore(FCore).CurrentGeneration; + for I:=0 to Snapshot.Count-1 do + begin + Entry:=TWSMessagePumpEntry(Snapshot[I]); + try + try + StartEntry(Entry,RunGeneration); + except + on E : Exception do + EntryEnded(Entry,E); + end; + finally + try + Entry.ReleaseReference; + except + on E : Exception do + ReportError(E); + end; + end; + end; + finally + Snapshot.Free; + end; + Exit; + end; + + Snapshot:=TList.Create; try - aList := List.LockList; + Entries:=FEntries.LockList; try - FReads:=[]; - for I := 0 to aList.Count - 1 do + for I:=0 to Entries.Count-1 do begin - aClient:= TWSClientConnection(aList.Items[I]); - if assigned(aClient.Transport) then - aClient.CheckIncoming(1); + Entry:=TWSMessagePumpEntry(Entries[I]); + Entry.AddReference; + Snapshot.Add(Entry); end; finally - List.UnlockList; + FEntries.UnlockList; end; - except - on E: Exception do - if Assigned(OnError) then - OnError(Self,E); + for I:=0 to Snapshot.Count-1 do + begin + Entry:=TWSMessagePumpEntry(Snapshot[I]); + try + if Entry.IsRegistered and Entry.Session.IsOpen then + try + IncomingResult:=Entry.Connection.CheckIncoming(0); + if IncomingResult=irClose then + EntryEnded(Entry,Nil); + except + on E : EWSReadInterrupted do + EntryEnded(Entry,Nil); + on E : Exception do + EntryEnded(Entry,E); + end; + finally + try + Entry.ReleaseReference; + except + on E : Exception do + ReportError(E); + end; + end; + end; + finally + Snapshot.Free; end; end; { TWSThreadMessagePump } +constructor TWSThreadMessagePump.Create(aOwner: TComponent); +begin + inherited Create(aOwner); + InitCriticalSection(FLifecycleLock); +end; + procedure TWSThreadMessagePump.Execute; +var + RunGeneration : LongInt; + DriverThread : TMessageDriverThread; begin - FThread:=TMessageDriverThread.Create(Self,@ThreadTerminated); + { Reap only a stopped generation. A harmless second Execute must never + interrupt the readers of the generation which is already running. } + RunGeneration:=TWSPumpCore(FCore).CurrentGeneration; + if not TWSPumpCore(FCore).IsRunning(RunGeneration) then + TryFinalize(0,True); + EnterCriticalSection(FLifecycleLock); + try + if Assigned(FThread) then + Exit; + if not TWSPumpCore(FCore).BeginRun(RunGeneration) then + Exit; + DriverThread:=TMessageDriverThread.Create(Self,Nil); + DriverThread.FRunGeneration:=RunGeneration; + FThread:=DriverThread; + try + DriverThread.Start; + except + FThread:=Nil; + DriverThread.Free; + TWSPumpCore(FCore).RequestStop; + raise; + end; + finally + LeaveCriticalSection(FLifecycleLock); + end; end; -procedure TWSThreadMessagePump.ThreadTerminated(Sender: TObject); +destructor TWSThreadMessagePump.Destroy; begin - FThread:=Nil; + if CurrentWSPumpCore=TWSPumpCore(FCore) then + Raise EWebSocketClient.Create('A websocket message pump cannot be freed from its own callback'); + RequestStop; + while not TryFinalize(100,True) do + ; + DoneCriticalSection(FLifecycleLock); + inherited Destroy; end; -procedure TWSThreadMessagePump.Terminate; +procedure TWSThreadMessagePump.RequestStop; var - lThread: TThread; - lCounter: Integer; + DriverThread : TThread; begin - lThread := FThread; - if Assigned(lThread) then - begin - lThread.Terminate; - - // Wait till it stops - lCounter := 0; - while Assigned(FThread) and (lCounter < 200) do // 5 second timeout - begin - Sleep(10); - Inc(lCounter); - end; + TWSPumpCore(FCore).RequestStop; + EnterCriticalSection(FLifecycleLock); + try + DriverThread:=FThread; + if Assigned(DriverThread) then + DriverThread.Terminate; + finally + LeaveCriticalSection(FLifecycleLock); + end; +end; - // If thread still hasn't finished, there's a serious problem - if Assigned(FThread) then - begin - FThread.OnTerminate:=Nil; - // Force cleanup as last resort - FThread := nil; +function TWSThreadMessagePump.TryFinalize(aTimeoutMs: Integer; + aInterruptReaders: Boolean): Boolean; +var + Started : QWord; + DriverThread : TThread; + DriverFinished : Boolean; + TimedOut : Boolean; +begin + Result:=False; + if CurrentWSPumpCore=TWSPumpCore(FCore) then + Exit; + Started:=TThread.GetTickCount64; + repeat + if aInterruptReaders then + InterruptConnections; + EnterCriticalSection(FLifecycleLock); + try + DriverThread:=FThread; + DriverFinished:=(not Assigned(DriverThread)) or DriverThread.Finished; + finally + LeaveCriticalSection(FLifecycleLock); end; + Result:=(TWSPumpCore(FCore).WorkerCount=0) and + DriverFinished; + if Result then + Break; + TimedOut:=(aTimeoutMs>=0) and + ((TThread.GetTickCount64-Started)>=QWord(aTimeoutMs)); + if TimedOut then + Exit(False); + if TThread.CurrentThread.ThreadID=MainThreadID then + CheckSynchronize(0); + TThread.Sleep(1); + until False; + + EnterCriticalSection(FLifecycleLock); + try + DriverThread:=FThread; + if Assigned(DriverThread) and DriverThread.Finished then + begin + DriverThread.WaitFor; + FThread:=Nil; + DriverThread.Free; + end; + finally + LeaveCriticalSection(FLifecycleLock); end; end; +procedure TWSThreadMessagePump.Terminate; +Const + MinStopGraceMs = 100; + MaxStopGraceMs = 1000; +Var + GraceMs : QWord; +begin + RequestStop; + { A callback running on one of this pump's reader threads may request a + stop, but it must not wait for itself. } + if CurrentWSPumpCore=TWSPumpCore(FCore) then + Exit; + + { Preserve the previous ability to stop and restart a healthy pump without + disconnecting its clients. Normally the bounded polling loop exits within + two intervals. Interrupt sockets only when that graceful stop fails. } + if Interval>0 then + GraceMs:=QWord(Interval)*2+10 + else + GraceMs:=MinStopGraceMs; + if GraceMsMaxStopGraceMs then + GraceMs:=MaxStopGraceMs; + + { Bounded finalization is essential when this call is itself executing in + a method which the reader synchronized to the calling thread. On timeout + the thread objects and core are retained; a later Terminate, Execute or + destructor reaps them after the callback unwinds. } + if TryFinalize(Integer(GraceMs),False) then + Exit; + { A worker which did not leave within the normal read timeout is inside an + exact read. Targeted, repeated interruption now forces only those + readers out; this second wait remains bounded for synchronized callers. } + TryFinalize(Integer(GraceMs),True); +end; + { TWSThreadMessagePump.TMessageDriverThread } -constructor TWSThreadMessagePump.TMessageDriverThread.Create(aPump: TWSThreadMessagePump; aTerminate : TNotifyEvent); +constructor TWSThreadMessagePump.TMessageDriverThread.Create( + aPump: TWSThreadMessagePump; aTerminate: TNotifyEvent); begin FPump:=aPump; OnTerminate:=aTerminate; - FreeOnTerminate:=True; - Inherited Create(False); + Inherited Create(True); + FreeOnTerminate:=False; end; procedure TWSThreadMessagePump.TMessageDriverThread.Execute; - +var + PreviousCore : TWSPumpCore; begin - While Not Terminated do - if FPump.CheckConnections then - FPump.ReadConnections - else + PreviousCore:=CurrentWSPumpCore; + CurrentWSPumpCore:=TWSPumpCore(FPump.FCore); + try + while (not Terminated) and + TWSPumpCore(FPump.FCore).IsRunning(FRunGeneration) do begin - TThread.Sleep(FPump.Interval); + try + if FPump.CheckConnections then + FPump.ReadConnections; + except + on E : Exception do + FPump.ReportError(E); end; - // OnTerminate is called in a synchronize. However, if no-one calls CheckSynchronize, it is never called. - // So we call it ourselves. - if assigned(OnTerminate) then - begin - OnTerminate(Self); - OnTerminate:=Nil; - end; + if (not Terminated) and + TWSPumpCore(FPump.FCore).IsRunning(FRunGeneration) then + if FPump.Interval>0 then + TThread.Sleep(FPump.Interval) + else + TThread.Sleep(1); + end; + finally + CurrentWSPumpCore:=PreviousCore; + end; end; end. diff --git a/packages/fcl-web/tests/wsshutdowntest.pas b/packages/fcl-web/tests/wsshutdowntest.pas new file mode 100644 index 00000000..bbdcceb4 --- /dev/null +++ b/packages/fcl-web/tests/wsshutdowntest.pas @@ -0,0 +1,4897 @@ +{ + Loopback regression test for the fpwebsocketclient message pump. + + Written to evaluate the bounded WebSocket client rewrite. The rewrite uses + one leased session per connection generation, stable pump registrations, + callback admission barriers and independent readers so one incomplete frame + cannot starve another connection. Shutdown interrupts only readers which do + not leave during the bounded graceful phase. + + On Windows, a stuck read is interrupted with raw socket shutdown followed + by CancelIoEx. The socket descriptor and TLS objects remain alive until the + reader thread has exited. + + Scenarios + --------- + 1 HTTP upgrade handshake against FPC's own TWebSocketServer, and echo. + 2 Repeated Execute/Terminate on a healthy pump; each echo is matched by + content so a stale reply cannot pass for a fresh one. + 3 Peer close: exactly one OnDisconnect, and the client goes inactive. + 4 Terminate while a healthy connection is idle; the connection remains + usable after restarting the pump. + 5 Upgrade, echo and peer close over TLS. + 6 Partial-frame stall over plain TCP. + 7 Partial-frame stall over TLS, with the reader blocked inside SSL_read. + 8 Terminate called from the OnDisconnect callback on the reader thread. + 9 A healthy client sharing one pump with a partial-frame stalled client. + 10 Terminate while a callback waits in TThread.Synchronize for the main + thread - the one wait that interrupting a socket cannot end. + 11 A connection that an earlier OnDisconnect callback destroys while the + pump still holds it in its notification queue. + 12 An exception from a later client while an earlier one is queued for + notification. + 13 Terminate called from a method the reader thread synchronized onto + the main thread. + 14 A close-frame control callback that disconnects an earlier client + while the pump is still walking its list. + 15 The peer closes, or resets, the connection in the middle of a frame. + 16 An OnError handler that synchronizes a method disconnecting another + client of the same pump. + 17 The same from OnMessage, where message callbacks run while the pump + holds its list lock. + 18 A client freed from the main thread while its own message callback + still runs; the connection must outlive the callback. + 19 A client freed from the main thread while its own callback waits in + TThread.Synchronize for that thread. + 20 A client freed by its owner while the pump's OnDisconnect + notification for it is still running. + 21 A method synchronized by a client's own OnMessage disconnects that + client. + 22 The same method disconnects and reconnects the client. + 23 OnDisconnect synchronizes a reconnect of the same client. + 24 OnDisconnect reconnects the same client directly on the reader. + 25 A client disconnects itself from its own OnMessage on the reader. + 26 The same from the control callback for a close frame. + 27 A client whose frame read stalls is freed, without Terminate. + 28 The same client is disconnected instead of freed. + 29 The owner disconnects and reconnects a client while the pump's + OnDisconnect notification for the peer close still runs. + 30 A worker thread frees a client whose callback waits in Synchronize. + 31 A reconnected client is freed while a callback of its old + connection still runs. + 32 The pump is freed while a connection its owner disconnected is + still inside a callback. + 33 A client disconnects itself from its OnMessage and OnDisconnect + reconnects it, inside the old connection's callback. + 34 The pump is reassigned while a callback runs on the old pump. + 35 The destructor of a released connection raises on the reader. + 36 OnDisconnect raises while an active client is freed. + 37 An incomplete frame on one client does not starve a healthy client + sharing the same pump. + + Scenarios 6, 7 and 9 deliberately put the reader inside a genuinely + blocking read. The other scenarios exercise bounded polling, connection + state, callbacks and compatibility. + + How it runs + ----------- + Without arguments, the program re-executes itself once per scenario using + --scenario N. The parent supervises each child with a timeout, so a hung + scenario is killed and reported as HUNG without stopping the remaining + tests. Each child also has an internal watchdog. + + The separate --pump-first mode demonstrates the case where a pump is + destroyed before a client that still references it. It is excluded from + the numbered run. Unpatched main crashes there (exit 217); v3 of the + patch clears the reference in a Notification override and passes. + + The separate --interrupt-race mode stages a read that completes exactly + while the pump is interrupting sockets. It is excluded from the numbered + run because it asks about a window a few instructions wide, which only a + deliberately delayed copy of the patch can be expected to show. + + Child exit codes: 0 passed, 1 failed, 2 skipped, 99 watchdog, + 98 self-test failure. + + Runner exit code: number of failed, hung and crashed scenarios. +} + +program wsshutdowntest; + +{$mode objfpc}{$H+} + +uses + {$IFDEF UNIX}cthreads, BaseUnix,{$ENDIF} + SysUtils, Classes, process, sockets, ssockets, sslbase, sslsockets, + opensslsockets, sha1, base64, + fpwebsocket, fpcustwsserver, fpwebsocketserver, fpwebsocketclient; + +Const + WaitLimitMs = 5000; // how long a scenario waits for an expected event + ScenarioLimitS = 45; // in-child watchdog budget per scenario + ChildLimitMs = 60000; // runner's hard limit per child process + PollMs = 5; + +{ --------------------------------------------------------------------- + Cross-thread flags and counters. + + Callbacks fire on the pump's reader thread and on server threads, so + nothing they touch may be a plain field read from the main thread. + Counters are LongInt manipulated only through InterLocked*; strings are + guarded by a critical section. + --------------------------------------------------------------------- } + +Function ReadCounter(Var aValue : LongInt) : LongInt; +begin + Result:=InterLockedExchangeAdd(aValue,0); +end; + +Procedure BumpCounter(Var aValue : LongInt); +begin + InterLockedIncrement(aValue); +end; + +{ --------------------------------------------------------------------- + Output and per-scenario bookkeeping + --------------------------------------------------------------------- } + +Var + OutLock : TRTLCriticalSection; + StateLock : TRTLCriticalSection; + ScenariosPassed : Integer = 0; + ScenariosFailed : Integer = 0; + ScenariosSkipped : Integer = 0; + CurrentName : String = ''; + CurrentFailed : Boolean = False; + CurrentDeadline : QWord = 0; // 0 = watchdog idle + +Procedure Say(Const aLine : String); +begin + EnterCriticalSection(OutLock); + try + Writeln(aLine); + Flush(Output); + finally + LeaveCriticalSection(OutLock); + end; +end; + +Procedure SetDeadline(Const aName : String; aDeadline : QWord); +begin + EnterCriticalSection(StateLock); + try + CurrentName:=aName; + CurrentDeadline:=aDeadline; + finally + LeaveCriticalSection(StateLock); + end; +end; + +Procedure GetDeadline(Out aName : String; Out aDeadline : QWord); +begin + EnterCriticalSection(StateLock); + try + aName:=CurrentName; + aDeadline:=CurrentDeadline; + finally + LeaveCriticalSection(StateLock); + end; +end; + +Procedure BeginScenario(Const aName : String); +begin + CurrentFailed:=False; + SetDeadline(aName,GetTickCount64+QWord(ScenarioLimitS)*1000); + Say('--- '+aName); +end; + +{ One assertion inside the current scenario. } +Procedure Check(Const aWhat : String; aOK : Boolean; Const aDetail : String = ''); +begin + if aOK then + Say(' ok '+aWhat) + else + begin + if aDetail='' then + Say(' FAIL '+aWhat) + else + Say(' FAIL '+aWhat+' ('+aDetail+')'); + CurrentFailed:=True; + end; +end; + +{ In a child process each of these ends the run, so the exit code carries + the single scenario's verdict back to the runner. } +Procedure EndScenario; +begin + SetDeadline('',0); + if CurrentFailed then + begin + Inc(ScenariosFailed); + Say(' => FAILED'); + end + else + begin + Inc(ScenariosPassed); + Say(' => passed'); + end; +end; + +Procedure SkipScenario(Const aName, aReason : String); +begin + SetDeadline('',0); + Inc(ScenariosSkipped); + Say('--- '+aName); + Say(' => SKIPPED: '+aReason); +end; + +{ --------------------------------------------------------------------- + Watchdog. A scenario that hangs would otherwise make the whole run + stall silently, and the hanging call can never report its own failure. + --------------------------------------------------------------------- } + +Type + TWatchdog = Class(TThread) + Public + Procedure Execute; override; + end; + +Procedure TWatchdog.Execute; +Var + D : QWord; + N : String; +begin + While not Terminated do + begin + GetDeadline(N,D); + if (D<>0) and (GetTickCount64>D) then + begin + Say(''); + Say('WATCHDOG: scenario "'+N+'" exceeded ' + +IntToStr(ScenarioLimitS)+' s - it is hung.'); + Say('This is itself the finding: the call under test never returned.'); + Flush(Output); + Halt(99); + end; + Sleep(100); + end; +end; + +{ --------------------------------------------------------------------- + Ports. Fixed ports make two concurrent runs interfere, so each run + picks its own base and walks upwards from there. + --------------------------------------------------------------------- } + +Var + PortSeq : Integer = 0; + PortBase : Integer; + +Function NextPort : Word; +begin + Inc(PortSeq); + Result:=PortBase+PortSeq; +end; + +{ Wait until a counter reaches at least aWanted, or the limit expires. } +Function WaitForCount(Var aCounter : LongInt; aWanted : LongInt; + aLimitMs : Integer = WaitLimitMs) : Boolean; +Var + Waited : Integer; +begin + Waited:=0; + While (ReadCounter(aCounter)=aWanted; +end; + +{ --------------------------------------------------------------------- + Sec-WebSocket-Accept, needed by the hand-rolled stall peer. + --------------------------------------------------------------------- } + +Function CalcAccept(Const aKey : String) : String; +Var + Hash : TSHA1Digest; + B : TBytes; +begin + Hash:=SHA1String(Trim(aKey)+SSecWebSocketGUID); + SetLength(B,SizeOf(Hash)); + Move(Hash,B[0],Length(B)); + Result:=EncodeBytesBase64(B); +end; + +{ RFC 6455 section 1.3 known-answer vector. If this is wrong every stall + scenario would fail for the wrong reason, so it is checked first. } +Function SelfTestAccept : Boolean; +Const + RFCKey = 'dGhlIHNhbXBsZSBub25jZQ=='; + RFCExpect = 's3pPLMBiTxaQ9kYGzzhZRbK+xOo='; +Var + Got : String; +begin + Got:=CalcAccept(RFCKey); + Result:=Got=RFCExpect; + if not Result then + Say(' FAIL RFC 6455 accept vector: got '+Got+', expected '+RFCExpect); +end; + +{ --------------------------------------------------------------------- + Echo server. Three behaviours: echo, close after the first message, or + send half a frame and then stay silent. + --------------------------------------------------------------------- } + +Type + TServerMode = (smEcho, smCloseAfterMessage, smStallAfterMessage, + smCloseFrameAfterMessage); + + TEchoServer = Class + Private + FServer : TWebSocketServer; + FPort : Word; + FMode : TServerMode; + FCert : TCertificateData; + FReceived : LongInt; + FHalfSent : LongInt; + FErrors : LongInt; + FRelease : LongInt; // set by the test to let a parked callback return + FStopping : LongInt; // set during teardown so a parked callback exits + FLastError : String; + FErrLock : TRTLCriticalSection; + Procedure DoMessage(Sender : TObject; Const aMessage : TWSMessage); + Procedure DoGetHandler(Sender : TObject; Const aUseSSL : Boolean; Out aHandler : TSocketHandler); + Procedure NoteError(Const aMsg : String); + Public + Constructor Create(aPort : Word; aMode : TServerMode; aUseSSL : Boolean); + Destructor Destroy; override; + Procedure Start; + Procedure Stop; + Function LastError : String; + Procedure Release; + Property Port : Word Read FPort; + Property Received : LongInt Read FReceived; + Property HalfSent : LongInt Read FHalfSent; + Property Errors : LongInt Read FErrors; + end; + +Constructor TEchoServer.Create(aPort : Word; aMode : TServerMode; aUseSSL : Boolean); +begin + InitCriticalSection(FErrLock); + FPort:=aPort; + FMode:=aMode; + FServer:=TWebSocketServer.Create(Nil); + FServer.Port:=aPort; + FServer.Host:='127.0.0.1'; + FServer.ThreadedAccept:=True; + FServer.ThreadMode:=wtmThread; + FServer.OnMessageReceived:=@DoMessage; + if aUseSSL then + begin + { TWebSocketServer.CertificateData is declared but never instantiated, + so the built-in CreateSSLSocketHandler dereferences nil. We supply + the handler ourselves and own the certificate. The client runs with + VerifyPeerCert:=False, so a generated self-signed cert suffices. } + FCert:=TCertificateData.Create; + FCert.HostName:='127.0.0.1'; + FServer.OnGetSocketHandler:=@DoGetHandler; + FServer.UseSSL:=True; + end; +end; + +Destructor TEchoServer.Destroy; +begin + Stop; + FreeAndNil(FServer); + FreeAndNil(FCert); + DoneCriticalSection(FErrLock); + inherited Destroy; +end; + +Procedure TEchoServer.NoteError(Const aMsg : String); +begin + EnterCriticalSection(FErrLock); + try + FLastError:=aMsg; + finally + LeaveCriticalSection(FErrLock); + end; + BumpCounter(FErrors); +end; + +Function TEchoServer.LastError : String; +begin + EnterCriticalSection(FErrLock); + try + Result:=FLastError; + finally + LeaveCriticalSection(FErrLock); + end; +end; + +Procedure TEchoServer.DoGetHandler(Sender : TObject; Const aUseSSL : Boolean; + Out aHandler : TSocketHandler); +Var + S : TSSLSocketHandler; + CK : TCertAndKey; +begin + aHandler:=Nil; + if not aUseSSL then + Exit; + S:=TSSLSocketHandler.GetDefaultHandler; + try + if FCert.NeedCertificateData then + begin + S.CertGenerator.HostName:=FCert.HostName; + CK:=S.CertGenerator.CreateCertificateAndKey; + FCert.Certificate.Value:=CK.Certificate; + FCert.PrivateKey.Value:=CK.PrivateKey; + end; + S.CertificateData:=FCert; + aHandler:=S; + except + On E : Exception do + begin + S.Free; + NoteError('handler: '+E.Message); + Raise; + end; + end; +end; + +Procedure TEchoServer.DoMessage(Sender : TObject; Const aMessage : TWSMessage); +Var + Con : TWSServerConnection; + Half : TBytes; +begin + BumpCounter(FReceived); + try + Con:=Sender as TWSServerConnection; + Case FMode of + smEcho: + Con.Send(aMessage.AsString); + smCloseAfterMessage: + Con.Disconnect; + smCloseFrameAfterMessage: + { A proper close frame, 1000 = normal closure. Disconnect above + only closes the socket, and then the client never dispatches a + close control event. } + Con.Close('bye',1000); + smStallAfterMessage: + begin + { A final text frame announcing five payload bytes that never + arrive. Written through the transport, so over TLS this lands + inside a TLS record and parks the client in SSL_read. } + Half:=[$81,$05]; + Con.Transport.WriteBuffer(Half); + BumpCounter(FHalfSent); + { Park here instead of returning. If this callback returned, the + server would resume reading, notice the FIN that the client's + own SHUT_RDWR produces, and close the connection - and that peer + activity, not the local shutdown, would release the client's + read. The plain-TCP peer holds its connection open in exactly + the same way, so both scenarios differ only in transport. } + While (ReadCounter(FRelease)=0) and (ReadCounter(FStopping)=0) do + Sleep(25); + end; + end; + except + On E : Exception do + NoteError('message: '+E.Message); + end; +end; + +Procedure TEchoServer.Release; +begin + InterLockedExchange(FRelease,1); +end; + +Procedure TEchoServer.Start; +begin + FServer.Active:=True; +end; + +Procedure TEchoServer.Stop; +begin + InterLockedExchange(FStopping,1); + if Assigned(FServer) and FServer.Active then + Try + FServer.Active:=False; + except + On E : Exception do + NoteError('stop: '+E.Message); + end; +end; + + +{ --------------------------------------------------------------------- + Read instrumentation. + + The stall scenarios need to prove that the reader is actually parked in + a payload read, not merely that time has passed. TWSConnection's + CheckIncoming is not virtual, but GetTransport is, so a delegating + IWSTransport can count entries and exits of the blocking reads. + + While ReadEntries > ReadExits the frame reader is inside a read. Each + scenario runs in its own process, so plain globals are safe here. + --------------------------------------------------------------------- } + +Var + ReadEntries : LongInt = 0; + ReadExits : LongInt = 0; + +Function ReaderIsInsideRead : Boolean; +begin + Result:=ReadCounter(ReadEntries)>ReadCounter(ReadExits); +end; + +{ True once the reader has stayed inside a read for a while. A single + sample can hit the short read of a frame header, which ends by itself + before the read that stalls begins. } +Function WaitReaderParked(aStableMs : Integer = 100) : Boolean; +Var + Waited, Inside : Integer; +begin + Waited:=0; + Inside:=0; + While (Inside=aStableMs; +end; + +Type + { IWSTransport is declared under {$INTERFACES CORBA}, so references + through it are NOT reference counted. The wrapper therefore has to be + owned and freed explicitly by the connection that installs it. } + TInstrumentedTransport = Class(TObject, IWSTransport) + Private + FInner : IWSTransport; + Public + Constructor Create(aInner : IWSTransport); + Function CanRead(aTimeOut: Integer) : Boolean; + Procedure ReadBuffer(aBytes : TBytes); + Function ReadBytes(var aBytes : TBytes; aCount : Integer) : Integer; + Function WriteBytes(aBytes : TBytes; aCount : Integer) : Integer; + Procedure WriteBuffer(aBytes : TBytes); + Function ReadLn : String; + Function PeerIP : string; + Function PeerPort : word; + end; + +Constructor TInstrumentedTransport.Create(aInner : IWSTransport); +begin + FInner:=aInner; +end; + +{ CanRead is a bounded select, not a blocking read - not counted. } +Function TInstrumentedTransport.CanRead(aTimeOut: Integer) : Boolean; +begin + Result:=FInner.CanRead(aTimeOut); +end; + +Procedure TInstrumentedTransport.ReadBuffer(aBytes : TBytes); +begin + BumpCounter(ReadEntries); + try + FInner.ReadBuffer(aBytes); + finally + BumpCounter(ReadExits); + end; +end; + +Function TInstrumentedTransport.ReadBytes(var aBytes : TBytes; aCount : Integer) : Integer; +begin + BumpCounter(ReadEntries); + try + Result:=FInner.ReadBytes(aBytes,aCount); + finally + BumpCounter(ReadExits); + end; +end; + +Function TInstrumentedTransport.WriteBytes(aBytes : TBytes; aCount : Integer) : Integer; +begin + Result:=FInner.WriteBytes(aBytes,aCount); +end; + +Procedure TInstrumentedTransport.WriteBuffer(aBytes : TBytes); +begin + FInner.WriteBuffer(aBytes); +end; + +Function TInstrumentedTransport.ReadLn : String; +begin + BumpCounter(ReadEntries); + try + Result:=FInner.ReadLn; + finally + BumpCounter(ReadExits); + end; +end; + +Function TInstrumentedTransport.PeerIP : string; +begin + Result:=FInner.PeerIP; +end; + +Function TInstrumentedTransport.PeerPort : word; +begin + Result:=FInner.PeerPort; +end; + +Type + TInstrumentedConnection = Class(TWebSocketClientConnection) + Private + FWrapper : TInstrumentedTransport; // owned, see the note above + Protected + Function GetTransport : IWSTransport; override; + Public + Destructor Destroy; override; + end; + +Function TInstrumentedConnection.GetTransport : IWSTransport; +begin + if FWrapper=Nil then + FWrapper:=TInstrumentedTransport.Create(inherited GetTransport); + Result:=FWrapper; +end; + +Destructor TInstrumentedConnection.Destroy; +begin + inherited Destroy; + FreeAndNil(FWrapper); +end; + +Type + TInstrumentedClient = Class(TWebsocketClient) + Protected + Function CreateClientConnection(aTransport : TWSClientTransport) : TWebsocketClientConnection; override; + end; + +Function TInstrumentedClient.CreateClientConnection(aTransport : TWSClientTransport) : TWebsocketClientConnection; +begin + Result:=TInstrumentedConnection.Create(Self,aTransport,Options); +end; + + +{ --------------------------------------------------------------------- + Lifetime probe. + + Scenario 11 destroys a connection that the pump has already collected for + notification. The question is whether the collected pointer is used + afterwards, so a probe connection notes its own address as its instance + is released, and DoDisconnect consults that note before anything else. + + FreeInstance deliberately keeps the block instead of returning it to the + heap. The destructor has run and CleanupInstance has finalised the + fields, so the object is destroyed in every sense that matters - but the + memory and the VMT pointer stay valid, so a later call lands in this + class and can be counted rather than jumping into whatever the heap + manager has since written there. The instances leak on purpose; the + scenario ends the process shortly afterwards. + --------------------------------------------------------------------- } + +Const + MaxDestroyed = 256; + +Var + DestroyedAddrs : Array[0..MaxDestroyed-1] of Pointer; + DestroyedCount : LongInt = 0; + DestroyLock : TRTLCriticalSection; + UseAfterFree : LongInt = 0; + ContinuedAfterFree : LongInt = 0; // a connection kept working after its destruction + ProbeClientDestroying : LongInt = 0; // Free of a probe client has begun + RaiseOnDestroyAddr : Pointer = Nil; // the probe connection whose destructor raises + DestructorRaised : LongInt = 0; // that destructor ran and raised + DestructorThread : TThreadID = 0; // the thread it ran on + WatchSlowDone : PLongInt = Nil; // the callback counter read by that destructor + CallbackDoneAtDestroy : LongInt = -1; + ProbePumpDestroying : LongInt = 0; // Free of a probe pump has begun + +{ Markers are published and read under a lock, so a reader never meets a + slot that is counted but not yet written. } +Procedure NoteDestroyed(aAddr : Pointer); +begin + EnterCriticalSection(DestroyLock); + try + if DestroyedCountNil) and (Pointer(Self)=RaiseOnDestroyAddr) then + begin + RaiseOnDestroyAddr:=Nil; + DestructorThread:=GetCurrentThreadID; + if WatchSlowDone<>Nil then + CallbackDoneAtDestroy:=InterLockedExchangeAdd(WatchSlowDone^,0); + BumpCounter(DestructorRaised); + Raise Exception.Create('deliberate failure in a connection destructor'); + end; +end; + +Type + { The handshake response a probe client creates. Keeps its memory like the + other probes, so a freed response can be recognised. } + TProbeResponse = Class(TWSHandShakeResponse) + Public + Procedure FreeInstance; override; + end; + + TProbeClient = Class(TWebsocketClient) + Protected + Function CreateClientConnection(aTransport : TWSClientTransport) : TWebsocketClientConnection; override; + Function CreateHandshakeResponse(aHeaders : TStrings) : TWSHandShakeResponse; override; + Public + Procedure FreeInstance; override; + Procedure BeforeDestruction; override; + end; + +Function TProbeClient.CreateClientConnection(aTransport : TWSClientTransport) : TWebsocketClientConnection; +begin + Result:=TProbeConnection.Create(Self,aTransport,Options); +end; + +Procedure TProbeClient.FreeInstance; +begin + { As for the probe connection: the memory stays, so a callback that + returns into a destroyed client component is counted, not a crash. } + CleanupInstance; + NoteDestroyed(Self); +end; + +Function TProbeClient.CreateHandshakeResponse(aHeaders : TStrings) : TWSHandShakeResponse; +begin + Result:=TProbeResponse.Create('',aHeaders); +end; + +Procedure TProbeClient.BeforeDestruction; +begin + { Marks the moment a thread has entered Free, before any teardown. } + BumpCounter(ProbeClientDestroying); + inherited BeforeDestruction; +end; + +Procedure TProbeResponse.FreeInstance; +begin + CleanupInstance; + NoteDestroyed(Self); +end; + +Type + { TWSMessagePump.List is protected, so asking the pump which connections + it still tracks needs a descendant rather than a cast. } + TProbePump = Class(TWSThreadMessagePump) + Public + Function ClientCount : Integer; + Function Tracks(aConnection : TWSClientConnection) : Boolean; + Procedure BeforeDestruction; override; + end; + +Function TProbePump.ClientCount : Integer; +Var + L : TList; +begin + L:=List.LockList; + try + Result:=L.Count; + finally + List.UnlockList; + end; +end; + +Function TProbePump.Tracks(aConnection : TWSClientConnection) : Boolean; +Var + L : TList; +begin + L:=List.LockList; + try + Result:=L.IndexOf(aConnection)>=0; + finally + List.UnlockList; + end; +end; + +Procedure TProbePump.BeforeDestruction; +begin + { Marks the moment a thread has entered the pump's Free. } + BumpCounter(ProbePumpDestroying); + inherited BeforeDestruction; +end; + +Type + { TWSErrorEvent is a method pointer, so the pump's OnError needs an + object to report into. } + TErrorSink = Class + Private + FErrors : LongInt; + FReadErrors : LongInt; + FLast : String; + FLock : TRTLCriticalSection; + Public + Constructor Create; + Destructor Destroy; override; + Procedure DoError(Sender : TObject; E : Exception); + Function LastError : String; + Property Errors : LongInt Read FErrors; + Property ReadErrors : LongInt Read FReadErrors; + end; + +Constructor TErrorSink.Create; +begin + InitCriticalSection(FLock); +end; + +Destructor TErrorSink.Destroy; +begin + DoneCriticalSection(FLock); + inherited Destroy; +end; + +Procedure TErrorSink.DoError(Sender : TObject; E : Exception); +begin + EnterCriticalSection(FLock); + try + FLast:=E.ClassName+': '+E.Message; + finally + LeaveCriticalSection(FLock); + end; + if E is EWSReadError then + BumpCounter(FReadErrors); + BumpCounter(FErrors); +end; + +Function TErrorSink.LastError : String; +begin + EnterCriticalSection(FLock); + try + Result:=FLast; + finally + LeaveCriticalSection(FLock); + end; +end; + + +{ --------------------------------------------------------------------- + Client wrapper. Counters rather than flags, so "exactly once" can be + asserted. + --------------------------------------------------------------------- } + +Type + TTestClient = Class + Private + FClient : TWebsocketClient; + FMessages : LongInt; + FDisconnects : LongInt; + FLast : String; + FLastLock : TRTLCriticalSection; + FTerminateOnDisconnect : TWSMessagePump; + FTerminateEntered : LongInt; // the callback reached the Terminate call + FTerminateReturned : LongInt; // Terminate returned normally + FTerminateRaised : LongInt; // Terminate raised an exception + FTerminateError : String; + FTermLock : TRTLCriticalSection; + FSyncOnMessage : Boolean; + FSyncEntered : LongInt; // the callback is about to call Synchronize + FSyncReturned : LongInt; // Synchronize returned + FRaiseOnMessage : Boolean; + FTearDownOnDisconnect : TCustomWebsocketClient; + FSyncTerminatePump : TWSMessagePump; + FTearDownOnCloseFrame : TCustomWebsocketClient; + FControlCloses : LongInt; + FSyncDisconnectOther : TCustomWebsocketClient; + FSyncOtherEntered : LongInt; // the synchronized method started + FSyncOtherReturned : LongInt; // the synchronized method returned + FSlowMessageMs : Integer; + FSlowEntered : LongInt; // the slow callback started + FSlowDone : LongInt; // the slow callback finished + FSlowSawDestroyed : LongInt; // its connection was destroyed before it finished + FSlowConnection : TObject; // the connection whose lifetime is checked + FSlowDisconnectMs : Integer; + FSlowDiscEntered : LongInt; // the slow OnDisconnect started + FSlowDiscDone : LongInt; // the slow OnDisconnect finished + FSlowDiscSawDestroyed : LongInt; + FSyncTermEntered : LongInt; // the synchronized method reached Terminate + FSyncTermReturned : LongInt; // Terminate returned inside it + FSyncTermRaised : LongInt; // Terminate raised inside it + FSlowComponent : TObject; // the client component whose lifetime is checked + FSlowSawComponentDestroyed : LongInt; + FSyncSelfAction : LongInt; // OnMessage synchronizes an action on this client + FDirectSelfAction : LongInt; // OnMessage runs it directly on the reader thread + FCloseFrameSelfDisconnect : Boolean; + FReconnectOnDisconnect : LongInt; // 1 directly in OnDisconnect, 2 synchronized from it + FPendingAction : LongInt; // 1 disconnect, 2 disconnect and connect, 3 connect + FSelfEntered : LongInt; // the action started + FSelfReturned : LongInt; // the action returned + FSelfRaised : LongInt; // the action raised + FSelfError : String; + FCallbackDone : LongInt; // OnMessage went on after the action + FOldResponse : TObject; // the handshake response of the connection in use + FOldResponseFreed : LongInt; // that response was freed while the callback still ran + FRaiseOnDisconnect : LongInt; // OnDisconnect raises once + FHoldUntilRelease : Boolean; // the slow callback waits for FReleaseHold instead of sleeping + FReleaseHold : LongInt; + FHoldTimedOut : LongInt; // the hold ended by its time limit, not by a release + Procedure DoNothing; + Procedure DoSyncTerminate; + Procedure DoControl(Sender : TObject; aType : TFrameType; Const aData : TBytes); + Procedure DoSyncDisconnectOther; + Procedure DoMessage(Sender : TObject; Const aMessage : TWSMessage); + Procedure DoDisconnect(Sender : TObject); + Procedure RunSelfAction; + Public + Constructor Create(aPort : Word; aPump : TWSMessagePump; aUseSSL : Boolean = False; + aInstrument : Boolean = False; aProbe : Boolean = False); + Destructor Destroy; override; + Function LastMessage : String; + Property Client : TWebsocketClient Read FClient; + Property Messages : LongInt Read FMessages; + Property Disconnects : LongInt Read FDisconnects; + { When set, OnDisconnect calls Terminate on this pump - from the reader + thread, which is what makes the patched Terminate join itself. } + Property TerminateOnDisconnect : TWSMessagePump + Read FTerminateOnDisconnect Write FTerminateOnDisconnect; + { When set, OnMessage parks in TThread.Synchronize. That call only + returns once someone runs CheckSynchronize on the main thread. } + Property SyncOnMessage : Boolean Read FSyncOnMessage Write FSyncOnMessage; + { When set, OnMessage raises an ordinary exception after counting the + message - an application callback that fails, nothing more. } + Property RaiseOnMessage : Boolean Read FRaiseOnMessage Write FRaiseOnMessage; + { When set, the first OnDisconnect disconnects this other client, which + frees its connection object. Ordinary application behaviour: one + connection drops, so the rest are torn down as well. } + Property TearDownOnDisconnect : TCustomWebsocketClient + Read FTearDownOnDisconnect Write FTearDownOnDisconnect; + { When set, OnMessage synchronizes a method that stops this pump - the + way an application hands a message to its main thread, which then + decides to close. } + Property SyncTerminatePump : TWSMessagePump + Read FSyncTerminatePump Write FSyncTerminatePump; + { When set, receiving a close frame disconnects this other client from + inside the control callback - that is, while the pump is still inside + CheckIncoming for the connection that received it. } + Property TearDownOnCloseFrame : TCustomWebsocketClient + Read FTearDownOnCloseFrame Write FTearDownOnCloseFrame; + { When set, OnMessage synchronizes a method that disconnects this other + client of the same pump - the main thread reacting to a message by + closing another connection. } + Property SyncDisconnectOther : TCustomWebsocketClient + Read FSyncDisconnectOther Write FSyncDisconnectOther; + { When set, OnMessage stays in the callback this long and then checks + whether the connection that delivered the message is still alive. } + Property SlowMessageMs : Integer Read FSlowMessageMs Write FSlowMessageMs; + Property SyncEntered : LongInt Read FSyncEntered; + Property SyncReturned : LongInt Read FSyncReturned; + Property TerminateEntered : LongInt Read FTerminateEntered; + Property TerminateReturned : LongInt Read FTerminateReturned; + Property TerminateRaised : LongInt Read FTerminateRaised; + Function TerminateError : String; + Function SelfError : String; + end; + +Constructor TTestClient.Create(aPort : Word; aPump : TWSMessagePump; aUseSSL : Boolean = False; + aInstrument : Boolean = False; aProbe : Boolean = False); +begin + InitCriticalSection(FLastLock); + InitCriticalSection(FTermLock); + if aProbe then + FClient:=TProbeClient.Create(Nil) + else if aInstrument then + FClient:=TInstrumentedClient.Create(Nil) + else + FClient:=TWebsocketClient.Create(Nil); + FClient.HostName:='127.0.0.1'; + FClient.Port:=aPort; + FClient.Resource:='/'; + FClient.UseSSL:=aUseSSL; + FClient.MessagePump:=aPump; + FClient.OnMessageReceived:=@DoMessage; + FClient.OnDisconnect:=@DoDisconnect; + FClient.OnControl:=@DoControl; +end; + +Destructor TTestClient.Destroy; +begin + FreeAndNil(FClient); + DoneCriticalSection(FTermLock); + DoneCriticalSection(FLastLock); + inherited Destroy; +end; + +Procedure TTestClient.DoNothing; +begin + { The body is irrelevant; what matters is that it runs on the main + thread, so the reader thread waits until the main thread services it. } +end; + +Procedure TTestClient.DoControl(Sender : TObject; aType : TFrameType; Const aData : TBytes); +Var + Other : TCustomWebsocketClient; +begin + if aType<>ftClose then + Exit; + BumpCounter(FControlCloses); + if Assigned(FTearDownOnCloseFrame) then + begin + Other:=FTearDownOnCloseFrame; + FTearDownOnCloseFrame:=Nil; + Other.Disconnect(False); + end; + if FCloseFrameSelfDisconnect then + begin + { Still inside CheckIncoming for this very connection, which sends the + close reply once the callback returns. } + FCloseFrameSelfDisconnect:=False; + InterLockedExchange(FPendingAction,1); + RunSelfAction; + end; +end; + +Procedure TTestClient.DoSyncDisconnectOther; +Var + Other : TCustomWebsocketClient; +begin + BumpCounter(FSyncOtherEntered); + Other:=FSyncDisconnectOther; + FSyncDisconnectOther:=Nil; + if Assigned(Other) then + Other.Disconnect(False); + BumpCounter(FSyncOtherReturned); +end; + +Procedure TTestClient.DoSyncTerminate; +begin + { Runs on the main thread, while the reader thread waits in Synchronize + for this very method to return. } + BumpCounter(FSyncTermEntered); + try + FSyncTerminatePump.Terminate; + BumpCounter(FSyncTermReturned); + except + On E : Exception do + BumpCounter(FSyncTermRaised); + end; +end; + +Procedure TTestClient.DoMessage(Sender : TObject; Const aMessage : TWSMessage); +Var + HoldDeadline : QWord; + Released : Boolean; +begin + EnterCriticalSection(FLastLock); + try + FLast:=aMessage.AsString; + finally + LeaveCriticalSection(FLastLock); + end; + BumpCounter(FMessages); + if FSyncOnMessage then + begin + BumpCounter(FSyncEntered); + TThread.Synchronize(TThread.CurrentThread,@DoNothing); + BumpCounter(FSyncReturned); + end; + if Assigned(FSyncTerminatePump) then + TThread.Synchronize(TThread.CurrentThread,@DoSyncTerminate); + if Assigned(FSyncDisconnectOther) then + TThread.Synchronize(TThread.CurrentThread,@DoSyncDisconnectOther); + if ReadCounter(FSyncSelfAction)<>0 then + begin + InterLockedExchange(FPendingAction,InterLockedExchange(FSyncSelfAction,0)); + TThread.Synchronize(TThread.CurrentThread,@RunSelfAction); + BumpCounter(FCallbackDone); + end; + if ReadCounter(FDirectSelfAction)<>0 then + begin + InterLockedExchange(FPendingAction,InterLockedExchange(FDirectSelfAction,0)); + RunSelfAction; + { Still inside the old connection's callback: whatever the connection + refers to has to be alive here. } + if Assigned(FOldResponse) and WasDestroyed(Pointer(FOldResponse)) then + BumpCounter(FOldResponseFreed); + BumpCounter(FCallbackDone); + end; + if FSlowMessageMs>0 then + begin + BumpCounter(FSlowEntered); + if FHoldUntilRelease then + begin + { Held until another thread observed the operation this callback is + meant to overlap; bounded, so a hang elsewhere stays visible. } + HoldDeadline:=TThread.GetTickCount64+10000; + Released:=False; + Repeat + { The reason the loop ends is taken inside it: a release arriving + just after the deadline must not hide the timeout. } + if ReadCounter(FReleaseHold)<>0 then + Released:=True + else if TThread.GetTickCount64>=HoldDeadline then + Break + else + Sleep(PollMs); + until Released; + if not Released then + BumpCounter(FHoldTimedOut); + end + else + Sleep(FSlowMessageMs); + { Sender is the client component here, not the connection, so the + connection's address was noted before the message was sent. Only the + pointer is looked up; the probe keeps a destroyed connection's memory, + so this is safe either way. } + if WasDestroyed(Pointer(FSlowConnection)) then + BumpCounter(FSlowSawDestroyed); + if Assigned(FSlowComponent) and WasDestroyed(Pointer(FSlowComponent)) then + BumpCounter(FSlowSawComponentDestroyed); + BumpCounter(FSlowDone); + end; + if FRaiseOnMessage then + Raise Exception.Create('deliberate failure in an OnMessage callback'); +end; + +Procedure TTestClient.DoDisconnect(Sender : TObject); +Var + Other : TCustomWebsocketClient; + Mode : LongInt; +begin + BumpCounter(FDisconnects); + if InterLockedExchange(FRaiseOnDisconnect,0)<>0 then + Raise Exception.Create('deliberate failure in an OnDisconnect handler'); + if FSlowDisconnectMs>0 then + begin + BumpCounter(FSlowDiscEntered); + Sleep(FSlowDisconnectMs); + if WasDestroyed(Pointer(FSlowConnection)) then + BumpCounter(FSlowDiscSawDestroyed); + BumpCounter(FSlowDiscDone); + end; + if ReadCounter(FReconnectOnDisconnect)<>0 then + begin + { Once only: the reconnected client gets its own OnDisconnect later. } + Mode:=InterLockedExchange(FReconnectOnDisconnect,0); + InterLockedExchange(FPendingAction,3); + if Mode=2 then + TThread.Synchronize(TThread.CurrentThread,@RunSelfAction) + else + RunSelfAction; + end; + if Assigned(FTearDownOnDisconnect) then + begin + { Once only: the teardown itself produces an OnDisconnect. } + Other:=FTearDownOnDisconnect; + FTearDownOnDisconnect:=Nil; + Other.Disconnect(False); + end; + if Assigned(FTerminateOnDisconnect) then + begin + { Three outcomes have to be told apart: a normal return, an exception, + and neither - the last one being a genuine self-join deadlock. On + glibc pthread_join detects self-join and returns EDEADLK, which FPC + ignores, so the failure surfaces later as an exception instead. } + BumpCounter(FTerminateEntered); + try + FTerminateOnDisconnect.Terminate; + BumpCounter(FTerminateReturned); + except + On E : Exception do + begin + EnterCriticalSection(FTermLock); + try + FTerminateError:=E.ClassName+': '+E.Message; + finally + LeaveCriticalSection(FTermLock); + end; + BumpCounter(FTerminateRaised); + end; + end; + end; +end; + +{ Runs the pending action on this client - on the reader thread or on the + main thread, depending on who calls it. } +Procedure TTestClient.RunSelfAction; +Var + Action : LongInt; + Tries : Integer; +begin + Action:=InterLockedExchange(FPendingAction,0); + if Action=0 then + Exit; + BumpCounter(FSelfEntered); + try + if Action in [1,2] then + FClient.Disconnect(False); + if Action in [2,3] then + begin + { The listener may refuse for a moment; the question here is not the + network, so retry like ConnectWithRetry does. } + Tries:=0; + repeat + try + FClient.Connect; + Tries:=-1; + except + On E : ESocketError do + begin + Inc(Tries); + if Tries>=20 then + Raise; + Say(Format(' reconnect refused, retry %d',[Tries])); + Sleep(50); + end; + end; + until Tries<0; + end; + BumpCounter(FSelfReturned); + except + On E : Exception do + begin + EnterCriticalSection(FTermLock); + try + FSelfError:=E.ClassName+': '+E.Message; + finally + LeaveCriticalSection(FTermLock); + end; + BumpCounter(FSelfRaised); + end; + end; +end; + +Function TTestClient.SelfError : String; +begin + EnterCriticalSection(FTermLock); + try + Result:=FSelfError; + finally + LeaveCriticalSection(FTermLock); + end; +end; + +Function TTestClient.TerminateError : String; +begin + EnterCriticalSection(FTermLock); + try + Result:=FTerminateError; + finally + LeaveCriticalSection(FTermLock); + end; +end; + +Function TTestClient.LastMessage : String; +begin + EnterCriticalSection(FLastLock); + try + Result:=FLast; + finally + LeaveCriticalSection(FLastLock); + end; +end; + +{ --------------------------------------------------------------------- + Plain-TCP stall peer: completes the upgrade by hand, then sends two + bytes of a frame header and holds the connection open. + --------------------------------------------------------------------- } + +{ SO_LINGER with a zero timeout turns the following close into a reset. + The struct differs between the two platforms, and the TLinger of the + sockets unit is the Unix shape, so it is spelled out here. } +Type + TAbortLinger = packed record + {$IFDEF WINDOWS} + l_onoff : Word; + l_linger : Word; + {$ELSE} + l_onoff : LongInt; + l_linger : LongInt; + {$ENDIF} + end; + +Function SetAbortiveClose(aSock : Longint) : Boolean; +Var + L : TAbortLinger; +begin + L.l_onoff:=1; + L.l_linger:=0; + Result:=fpsetsockopt(aSock,SOL_SOCKET,SO_LINGER,@L,SizeOf(L))=0; +end; + +Type + TStallServer = Class(TThread) + Private + FPort : Word; + FListener : TInetServer; + FHandshakes : LongInt; + FHalfSent : LongInt; + FRestSent : LongInt; + FSendRest : LongInt; + FAfterHalf : LongInt; + FClosedAfterHalf : LongInt; + FErrors : LongInt; + FLastError : String; + FErrLock : TRTLCriticalSection; + Procedure DoConnect(Sender : TObject; Data : TSocketStream); + Procedure NoteError(Const aMsg : String); + Public + Constructor Create(aPort : Word); + Destructor Destroy; override; + Procedure Execute; override; + Procedure Shutdown; + { Release the five payload bytes the half frame announced. Used to make + a parked read complete at a chosen moment. } + Procedure SendRest; + Function LastError : String; + Property Handshakes : LongInt Read FHandshakes; + Property HalfSent : LongInt Read FHalfSent; + Property RestSent : LongInt Read FRestSent; + { What happens after the half frame: 0 holds the connection, 1 closes it + gracefully, 2 resets it. Set before the client connects. } + Property AfterHalf : LongInt Read FAfterHalf Write FAfterHalf; + Property ClosedAfterHalf : LongInt Read FClosedAfterHalf; + Property Errors : LongInt Read FErrors; + end; + +Constructor TStallServer.Create(aPort : Word); +begin + InitCriticalSection(FErrLock); + FPort:=aPort; + FListener:=TInetServer.Create('127.0.0.1',aPort); + FListener.OnConnect:=@DoConnect; + FListener.QueueSize:=5; + FreeOnTerminate:=False; + Inherited Create(False); +end; + +{ Terminate, stop accepting, join - in that order - then release the + listener. The destructor must not be the thing that stops the thread. } +Procedure TStallServer.Shutdown; +begin + Terminate; + if Assigned(FListener) then + Try + FListener.StopAccepting(True); + except + // the accept loop may already be gone + end; + WaitFor; +end; + +Destructor TStallServer.Destroy; +begin + { Safe even if the caller already did it: Terminate and StopAccepting + are idempotent, and WaitFor on a finished thread returns at once. } + Shutdown; + FreeAndNil(FListener); + inherited Destroy; + DoneCriticalSection(FErrLock); +end; + +Procedure TStallServer.NoteError(Const aMsg : String); +begin + EnterCriticalSection(FErrLock); + try + FLastError:=aMsg; + finally + LeaveCriticalSection(FErrLock); + end; + BumpCounter(FErrors); +end; + +Function TStallServer.LastError : String; +begin + EnterCriticalSection(FErrLock); + try + Result:=FLastError; + finally + LeaveCriticalSection(FErrLock); + end; +end; + +Procedure TStallServer.DoConnect(Sender : TObject; Data : TSocketStream); + + { Read the request headers. Bounded, and gives up if the peer stops + talking, so shutting the test down cannot hang in here forever. } + Function ReadHeaders(Out aHeaders : String) : Boolean; + Var + C : Char; + Res : String; + N : Integer; + Idle : Integer; + begin + Res:=''; + Idle:=0; + While (Pos(#13#10#13#10,Res)=0) and (Length(Res)<8192) and (not Terminated) do + begin + { Data.Read goes straight into a blocking recv, so ask first. Without + this the loop could never observe Terminated or its own deadline, + and a stuck setup would look like the shutdown hang under test. } + if not Data.CanRead(PollMs*10) then + begin + Inc(Idle,PollMs*10); + if Idle>WaitLimitMs then + Break; + Continue; + end; + N:=Data.Read(C,1); + if N=1 then + begin + Res:=Res+C; + Idle:=0; + end + else + Break; // peer closed or errored + end; + aHeaders:=Res; + Result:=Pos(#13#10#13#10,Res)>0; + end; + + Function ExtractKey(Const aHeaders : String; Out aKey : String) : Boolean; + Var + L : TStringList; + I : Integer; + N : String; + begin + aKey:=''; + L:=TStringList.Create; + try + L.Text:=aHeaders; + For I:=0 to L.Count-1 do + begin + N:=L[I]; + if SameText(Copy(N,1,Length(SSecWebsocketKey)+1),SSecWebsocketKey+':') then + begin + aKey:=Trim(Copy(N,Length(SSecWebsocketKey)+2,Length(N))); + Break; + end; + end; + finally + L.Free; + end; + Result:=aKey<>''; + end; + + { Write everything or report failure - a short write would leave the + client waiting for bytes we never sent, which would look like the + stall we are trying to stage on purpose. } + Function WriteAll(Const aBuf; aCount : Integer) : Boolean; + Var + Written, N : Integer; + P : PByte; + begin + P:=@aBuf; + Written:=0; + While Written0 then + begin + Sleep(300); + if (ReadCounter(FAfterHalf)=2) and not SetAbortiveClose(Data.Handle) then + NoteError('SO_LINGER could not be set, the close will not be a reset') + else + BumpCounter(FClosedAfterHalf); + Exit; + end; + + { Hold the connection open. The client is now stuck waiting for a + payload; a shutdown, a close, or SendRest can free it. } + While not Terminated do + begin + if InterLockedExchange(FSendRest,0)=1 then + begin + Rest[0]:=$68; Rest[1]:=$65; Rest[2]:=$6C; { 'hel' } + Rest[3]:=$6C; Rest[4]:=$6F; { 'lo' } + if not WriteAll(Rest[0],5) then + NoteError('short write on the frame payload') + else + BumpCounter(FRestSent); + end; + Sleep(1); + end; + except + On E : Exception do + NoteError('connection: '+E.Message); + end; + finally + Data.Free; + end; +end; + +Procedure TStallServer.SendRest; +begin + InterLockedExchange(FSendRest,1); +end; + +Procedure TStallServer.Execute; +begin + try + FListener.StartAccepting; + except + On E : Exception do + if not Terminated then + NoteError('accept: '+E.Message); + end; +end; + +{ --------------------------------------------------------------------- + Is there a usable OpenSSL on this machine? + --------------------------------------------------------------------- } + +Var + TLSChecked : Boolean = False; + TLSAvailable : Boolean = False; + TLSReason : String = ''; + +Function HaveTLS : Boolean; +Var + H : TSSLSocketHandler; +begin + if not TLSChecked then + begin + TLSChecked:=True; + try + H:=TSSLSocketHandler.GetDefaultHandler; + try + TLSAvailable:=Assigned(H); + if not TLSAvailable then + TLSReason:='no default SSL handler registered'; + finally + H.Free; + end; + except + On E : Exception do + begin + TLSAvailable:=False; + TLSReason:=E.Message; + end; + end; + end; + Result:=TLSAvailable; +end; + +{ --------------------------------------------------------------------- + Scenario 1: upgrade handshake and echo + --------------------------------------------------------------------- } + +Procedure TestUpgradeAndEcho; +Var + Srv : TEchoServer; + Pump : TWSThreadMessagePump; + Cli : TTestClient; +begin + BeginScenario('upgrade handshake and echo'); + Srv:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TWSThreadMessagePump.Create(Nil); + Cli:=Nil; + try + Srv.Start; + Sleep(150); + Pump.Execute; + Cli:=TTestClient.Create(Srv.Port,Pump); + Cli.Client.Connect; + Check('client is active after connect',Cli.Client.Active); + Cli.Client.SendMessage('ping'); + Check('echo received',WaitForCount(Cli.FMessages,1)); + Check('echo content matches',Cli.LastMessage='ping','got "'+Cli.LastMessage+'"'); + Check('server saw the message',ReadCounter(Srv.FReceived)=1, + 'received='+IntToStr(ReadCounter(Srv.FReceived))); + Check('no server errors',ReadCounter(Srv.FErrors)=0,Srv.LastError); + Pump.Terminate; + finally + FreeAndNil(Cli); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 2: repeated Execute/Terminate on a healthy connection + --------------------------------------------------------------------- } + +Procedure TestRepeatedExecuteTerminate; +Var + Srv : TEchoServer; + Pump : TWSThreadMessagePump; + Cli : TTestClient; + I : Integer; + Expect : String; + Got : Boolean; +begin + BeginScenario('repeated Execute/Terminate keeps the connection usable'); + Srv:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TWSThreadMessagePump.Create(Nil); + Cli:=Nil; + try + Srv.Start; + Sleep(150); + Pump.Execute; + Cli:=TTestClient.Create(Srv.Port,Pump); + Cli.Client.Connect; + + For I:=1 to 5 do + begin + Pump.Terminate; + Pump.Execute; + Expect:='round'+IntToStr(I); + Cli.Client.SendMessage(Expect); + { Match on content: waiting for a count alone could be satisfied by + a reply from an earlier round. } + Got:=WaitForCount(Cli.FMessages,I,2000) and (Cli.LastMessage=Expect); + Check('round '+IntToStr(I)+' echoed correctly',Got, + 'last="'+Cli.LastMessage+'" expected="'+Expect+'"'); + if not Got then + Break; + end; + Check('client still active',Cli.Client.Active); + Pump.Terminate; + finally + FreeAndNil(Cli); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 3: peer close reaches the owner exactly once + --------------------------------------------------------------------- } + +Procedure TestPeerClose; +Var + Srv : TEchoServer; + Pump : TWSThreadMessagePump; + Cli : TTestClient; +begin + BeginScenario('peer close notifies the owner and clears Active'); + Srv:=TEchoServer.Create(NextPort,smCloseAfterMessage,False); + Pump:=TWSThreadMessagePump.Create(Nil); + Cli:=Nil; + try + Srv.Start; + Sleep(150); + Pump.Execute; + Cli:=TTestClient.Create(Srv.Port,Pump); + Cli.Client.Connect; + Cli.Client.SendMessage('bye'); + + Check('OnDisconnect fired',WaitForCount(Cli.FDisconnects,1)); + Check('client is no longer active',not Cli.Client.Active, + 'Active='+BoolToStr(Cli.Client.Active,True)); + { Give any duplicate notification time to show up before counting. } + Sleep(300); + Check('OnDisconnect fired exactly once',ReadCounter(Cli.FDisconnects)=1, + 'count='+IntToStr(ReadCounter(Cli.FDisconnects))); + Pump.Terminate; + finally + FreeAndNil(Cli); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 4: Terminate on an idle but healthy connection. + + The reader is in its bounded select here, not in a blocking payload + read - that case is scenarios 7 and 8. What is asserted is the + compatibility promise: stopping a healthy pump must be quick and must + leave the connection usable. + --------------------------------------------------------------------- } + +Procedure TestTerminateWhileIdle; +Var + Srv : TEchoServer; + Pump : TWSThreadMessagePump; + Cli : TTestClient; + Started, Elapsed : QWord; +begin + BeginScenario('Terminate on an idle healthy connection'); + Srv:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TWSThreadMessagePump.Create(Nil); + Cli:=Nil; + try + Srv.Start; + Sleep(150); + Pump.Execute; + Cli:=TTestClient.Create(Srv.Port,Pump); + Cli.Client.Connect; + Sleep(300); // let the reader settle into its polling loop + + { The client stays alive and registered across the Terminate. } + Started:=GetTickCount64; + Pump.Terminate; + Elapsed:=GetTickCount64-Started; + Check('Terminate returned promptly',Elapsed0 then + Say(' Terminate raised inside the callback: '+Cli.TerminateError) + else if ReadCounter(Cli.FTerminateReturned)>0 then + Say(' Terminate returned normally inside the callback') + else + Say(' Terminate was entered but neither returned nor raised - ' + +'the reader thread is waiting for itself'); + + Check('Terminate from the disconnect callback returns normally', + ReadCounter(Cli.FTerminateReturned)>0, + 'entered=1 raised='+IntToStr(ReadCounter(Cli.FTerminateRaised)) + +' error="'+Cli.TerminateError+'"'); + end; + finally + Cli.TerminateOnDisconnect:=Nil; + FreeAndNil(Cli); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 10: collateral damage of the forced interrupt. + + InterruptConnections shuts down every registered socket, not just the + one that is stuck. A healthy client sharing the pump is shut down too. + The patch's own comment promises that stopping a healthy pump does not + disconnect its clients, so this checks whether the healthy client + survives a Terminate that had to force its way out. + --------------------------------------------------------------------- } + +Procedure TestCollateralInterrupt; +Var + EchoSrv : TEchoServer; + StallSrv : TStallServer; + Pump : TWSThreadMessagePump; + Healthy, Stalled : TTestClient; + StallPort : Word; + SendErr : String; + ActiveAfter, Usable : Boolean; + DiscAfter : LongInt; + Waited : Integer; +begin + BeginScenario('healthy client sharing a pump with a stalled one'); + EchoSrv:=TEchoServer.Create(NextPort,smEcho,False); + StallPort:=NextPort; + StallSrv:=TStallServer.Create(StallPort); + Pump:=TWSThreadMessagePump.Create(Nil); + Healthy:=Nil; + Stalled:=Nil; + try + EchoSrv.Start; + Sleep(200); + Pump.Execute; + + Healthy:=TTestClient.Create(EchoSrv.Port,Pump); + Healthy.Client.Connect; + Healthy.Client.SendMessage('warmup'); + Check('healthy client works before the stall', + WaitForCount(Healthy.FMessages,1) and (Healthy.LastMessage='warmup')); + + Stalled:=TTestClient.Create(StallPort,Pump,False,True); + Stalled.Client.Connect; + Check('stall peer sent its partial frame',WaitForCount(StallSrv.FHalfSent,1), + StallSrv.LastError); + Waited:=0; + While (not ReaderIsInsideRead) and (Waited'' then + Check('healthy client still usable after the forced interrupt',False, + 'send raised: '+SendErr) + else + Check('healthy client still usable after the forced interrupt',Usable, + 'last="'+Healthy.LastMessage+'"'); + + { The state question, judged on the snapshot taken before the restart: + if the connection was broken, the owner should have been told. } + Check('client state and notification agree after the interrupt', + Usable or (not ActiveAfter) or (DiscAfter>0), + Format('unusable but Active=%s with %d disconnect(s) - the owner ' + +'was not told',[BoolToStr(ActiveAfter,True),DiscAfter])); + Pump.Terminate; + finally + FreeAndNil(Healthy); + FreeAndNil(Stalled); + FreeAndNil(Pump); + if Assigned(StallSrv) then + begin + StallSrv.Shutdown; + FreeAndNil(StallSrv); + end; + FreeAndNil(EchoSrv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 5: pump destroyed before the client that points at it. + + TCustomWebsocketClient.Destroy calls Disconnect, which calls + MessagePump.RemoveClient. SetMessagePump does register a + FreeNotification, but the class has no Notification override that + clears FMessagePump, so the pointer is never cleared. + + Unpatched main crashes here, so it is opt-in: --pump-first runs it alone. + --------------------------------------------------------------------- } + +Procedure TestPumpFreedBeforeClient; +Var + Srv : TEchoServer; + Pump : TWSThreadMessagePump; + Cli : TTestClient; +begin + BeginScenario('pump freed before an active client'); + Srv:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TWSThreadMessagePump.Create(Nil); + Cli:=Nil; + try + Srv.Start; + Sleep(150); + Pump.Execute; + Cli:=TTestClient.Create(Srv.Port,Pump); + Cli.Client.Connect; + Check('client is active',Cli.Client.Active); + Sleep(300); + + Say(' freeing the pump first (an access violation here is the finding)'); + FreeAndNil(Pump); + FreeAndNil(Cli); + Check('survived pump-before-client teardown',True); + finally + FreeAndNil(Cli); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + + +{ --------------------------------------------------------------------- + Probe mode --interrupt-race: an interrupt landing on a read that is just + completing. + + Not part of the numbered run, because it asks a question about a window + rather than about the patch's normal behaviour. + + Written against v2 of the patch, where InterruptRead claimed a transport + inside a read and published the request afterwards, and EndRead cleared + the state and read the request afterwards: between claim and publication + EndRead could see no request and return normally while the socket was + shut down a moment later. v3 claims and publishes in one atomic state + change, so against v3 and later this mode checks that a read completing + exactly while the pump interrupts is still reported. Running it against a + copy with a deliberate delay inside InterruptRead widens the window. + + Terminate waits gracefully for 100 ms and only then starts interrupting, + so the peer is armed to deliver the missing payload 120 ms after + Terminate is entered - about 20 ms into the interrupt loop. + --------------------------------------------------------------------- } + +Type + { The main thread is inside Terminate when the payload has to arrive, so + the peer is released from a thread of its own. } + TArmThread = Class(TThread) + Private + FServer : TStallServer; + FDelayMs : Integer; + Public + Constructor Create(aServer : TStallServer; aDelayMs : Integer); + Procedure Execute; override; + end; + +Constructor TArmThread.Create(aServer : TStallServer; aDelayMs : Integer); +begin + FServer:=aServer; + FDelayMs:=aDelayMs; + FreeOnTerminate:=False; + Inherited Create(False); +end; + +Procedure TArmThread.Execute; +begin + Sleep(FDelayMs); + FServer.SendRest; +end; + +Procedure TestInterruptOnCompletingRead; +Var + Srv : TStallServer; + Pump : TProbePump; + Cli : TTestClient; + Arm : TArmThread; + Con : TWSClientConnection; + Port : Word; + Waited : Integer; + Active, Tracked : Boolean; + Disc : LongInt; + SendErr : String; + Started, Elapsed : QWord; +begin + BeginScenario('interrupt landing on a read that is just completing'); + Port:=NextPort; + Srv:=TStallServer.Create(Port); + Pump:=TProbePump.Create(Nil); + Cli:=Nil; + Arm:=Nil; + try + Sleep(250); { let the peer's listener come up before connecting } + Pump.Execute; + Cli:=TTestClient.Create(Port,Pump,False,True); + Cli.Client.Connect; + Check('the peer sent its partial frame',WaitForCount(Srv.FHalfSent,1), + Srv.LastError); + Waited:=0; + While (not ReaderIsInsideRead) and (Waited=1) or (not Active), + 'the socket is gone, the client still reports Active, and no ' + +'disconnect was reported'); + Check('a connection the pump still tracks has a usable socket', + (not Tracked) or (SendErr=''), + 'it is still registered with a socket that has been shut down'); + finally + if Assigned(Arm) then + begin + Arm.WaitFor; + FreeAndNil(Arm); + end; + FreeAndNil(Cli); + FreeAndNil(Pump); + if Assigned(Srv) then + begin + Srv.Shutdown; + FreeAndNil(Srv); + end; + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario table, child entry point and runner + --------------------------------------------------------------------- } + +Type + TScenarioProc = Procedure; + +{ --------------------------------------------------------------------- + Scenario 10: a callback that waits for the main thread, while the main + thread is inside Terminate. + + TThread.Synchronize parks the reader until someone runs CheckSynchronize + on the main thread. Terminate polls on the main thread and does not, so + the two wait for each other. Interrupting the socket cannot help here: + the reader is not in a read. + + What differs between the versions is how they get out of it. A bounded + wait gives up and abandons the thread - unsafe, but it returns. An + unbounded one does not return at all. + --------------------------------------------------------------------- } +Procedure TestSynchronizeDuringTerminate; +Var + Srv : TEchoServer; + Pump : TWSThreadMessagePump; + Cli : TTestClient; + Started, Elapsed : QWord; + Reached : Boolean; +begin + BeginScenario('Terminate while a callback waits in Synchronize'); + Srv:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TWSThreadMessagePump.Create(Nil); + Cli:=Nil; + try + Srv.Start; + Sleep(150); + Pump.Execute; + Cli:=TTestClient.Create(Srv.Port,Pump); + Cli.SyncOnMessage:=True; + Cli.Client.Connect; + Cli.Client.SendMessage('sync'); + + Reached:=WaitForCount(Cli.FSyncEntered,1); + Check('the callback reached Synchronize',Reached); + if Reached then + begin + { Nobody has serviced it, so it must still be waiting. Without this the + scenario could pass with the reader long since finished. } + Sleep(200); + Check('the callback is still parked in Synchronize', + ReadCounter(Cli.FSyncReturned)=0, + 'it returned on its own, so the main thread is not the only one ' + +'that can service it and this scenario proves nothing'); + + Started:=TThread.GetTickCount64; + Pump.Terminate; + Elapsed:=TThread.GetTickCount64-Started; + Say(Format(' Terminate returned after %d ms',[Elapsed])); + Check('Terminate returns while a callback waits for the main thread',True); + + { Release the parked callback so that the teardown below is not itself + the thing that hangs. } + CheckSynchronize(0); + Say(' after CheckSynchronize the callback had ' + +BoolToStr(ReadCounter(Cli.FSyncReturned)>0,'returned','not returned')); + end; + finally + CheckSynchronize(0); + FreeAndNil(Cli); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + + +{ --------------------------------------------------------------------- + Scenario 11: a connection destroyed while the pump still has it queued. + + ReadConnections removes closed connections from both registries under the + list lock, collects them in a local list, and calls Disconnect on each of + them after unlocking. Between the first and the second of those calls the + application's OnDisconnect callback runs - and tearing the remaining + connections down from there is an ordinary thing for an application to do. + + Both peers close before the pump is started, so its very first pass finds + both sockets ready and collects both in the same pass. The first callback + then disconnects the second client, which frees its connection object, + and the pump goes on to call Disconnect on that pointer. + + A control round runs the same arrangement without the teardown, so a + finding below cannot be blamed on the setup. + --------------------------------------------------------------------- } + +Type + TQueuedRound = Record + DiscA, DiscB, Used : LongInt; + Armed : Boolean; + end; + +Function RunQueuedRound(aTearDown : Boolean) : TQueuedRound; +Var + SrvA, SrvB : TEchoServer; + Pump : TWSThreadMessagePump; + A, B : TTestClient; + Before : LongInt; +begin + Result.DiscA:=0; + Result.DiscB:=0; + Result.Used:=0; + Result.Armed:=False; + SrvA:=TEchoServer.Create(NextPort,smCloseAfterMessage,False); + SrvB:=TEchoServer.Create(NextPort,smCloseAfterMessage,False); + Pump:=TWSThreadMessagePump.Create(Nil); + A:=Nil; + B:=Nil; + try + SrvA.Start; + SrvB.Start; + Sleep(200); + + { The pump is deliberately left stopped. Both peers have to have closed + before its first pass, otherwise the two connections are collected in + two passes and the queue never holds more than one. } + A:=TTestClient.Create(SrvA.Port,Pump,False,False,True); + B:=TTestClient.Create(SrvB.Port,Pump,False,False,True); + A.Client.Connect; + B.Client.Connect; + A.Client.SendMessage('bye'); + B.Client.SendMessage('bye'); + Sleep(500); + + { Both peers must have closed before the pump's first pass. Otherwise the + two connections are collected one per pass, B is never queued while A's + callback runs, and a clean result below would mean nothing was tried. + CheckIncoming with DoRead=False answers "is there something waiting" + without consuming it. } + Result.Armed:=A.Client.Active and B.Client.Active + and (A.Client.Connection.CheckIncoming(50,False)=irWaiting) + and (B.Client.Connection.CheckIncoming(50,False)=irWaiting); + + if aTearDown then + A.TearDownOnDisconnect:=B.Client; + + Before:=ReadCounter(UseAfterFree); + Pump.Execute; + WaitForCount(A.FDisconnects,1); + Sleep(400); + Pump.Terminate; + + Result.DiscA:=ReadCounter(A.FDisconnects); + Result.DiscB:=ReadCounter(B.FDisconnects); + Result.Used:=ReadCounter(UseAfterFree)-Before; + finally + FreeAndNil(A); + FreeAndNil(B); + FreeAndNil(Pump); + FreeAndNil(SrvA); + FreeAndNil(SrvB); + end; +end; + +Procedure TestFreeQueuedDisconnect; +Var + Ctl, Prov : TQueuedRound; +begin + BeginScenario('a queued connection destroyed by an earlier callback'); + + Say(' control round: the callback does not tear anything down'); + Ctl:=RunQueuedRound(False); + Say(Format(' control: OnDisconnect A=%d B=%d, calls after destruction=%d', + [Ctl.DiscA,Ctl.DiscB,Ctl.Used])); + Check('control: both peers had closed before the pump started',Ctl.Armed, + 'they were not both waiting, so they were not collected in one pass'); + Check('control: both peer closes are reported exactly once', + (Ctl.DiscA=1) and (Ctl.DiscB=1), + Format('A=%d B=%d',[Ctl.DiscA,Ctl.DiscB])); + Check('control: nothing is used after destruction',Ctl.Used=0, + Format('%d call(s)',[Ctl.Used])); + + Say(' provocation round: the first callback disconnects the second client'); + Prov:=RunQueuedRound(True); + Say(Format(' provocation: OnDisconnect A=%d B=%d, calls after destruction=%d', + [Prov.DiscA,Prov.DiscB,Prov.Used])); + + { Two ways this round can decide nothing: the two peers were not collected + in the same pass, or the first callback never fired. Either has to be + reported as such rather than as a clean result. } + if not Prov.Armed then + Check('the provocation actually ran',False, + 'the two peers had not both closed before the pump started, so ' + +'they were not collected in one pass and B was never queued') + else if Prov.DiscA=0 then + Check('the provocation actually ran',False, + 'the first callback never fired, so the second client was never ' + +'torn down and this round decides nothing') + else + Check('the pump does not touch a connection destroyed by a callback', + Prov.Used=0, + Format('%d call(s) reached a connection whose destructor had ' + +'completed; the probe keeps that memory alive on purpose, ' + +'an ordinary build hands it back to the heap manager', + [Prov.Used])); + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 12: an exception in a later client, and the notifications that + were already pending. + + ReadConnections removes closed connections from both registries first and + notifies their owners afterwards, in a second loop. Both loops sit inside + one try..except. An ordinary exception from a *later* client - an + application callback that fails - therefore jumps past the notification + loop. The connections removed before it are then in neither registry and + have not been told, so nothing will ever look at them again. + + Two clients: the first one's peer closes, the second one's callback + raises. A control round with the same arrangement, minus the failing + callback, shows what the pass does when nothing interferes. + --------------------------------------------------------------------- } + +Type + TSkipRound = Record + Errors, DiscA : LongInt; + TrackedA, TrackedB, Armed : Boolean; + Error : String; + end; + +Function RunSkipRound(aRaise : Boolean) : TSkipRound; +Var + SrvA, SrvB : TEchoServer; + Pump : TProbePump; + Sink : TErrorSink; + A, B : TTestClient; + ConA, ConB : TWSClientConnection; +begin + Result.Errors:=0; + Result.DiscA:=0; + Result.TrackedA:=False; + Result.TrackedB:=False; + Result.Armed:=False; + Result.Error:=''; + SrvA:=TEchoServer.Create(NextPort,smCloseAfterMessage,False); + SrvB:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TProbePump.Create(Nil); + Sink:=TErrorSink.Create; + A:=Nil; + B:=Nil; + try + Pump.OnError:=@Sink.DoError; + SrvA.Start; + SrvB.Start; + Sleep(200); + + { Connect order is list order, so A is visited before B. } + A:=TTestClient.Create(SrvA.Port,Pump); + B:=TTestClient.Create(SrvB.Port,Pump); + A.Client.Connect; + B.Client.Connect; + B.RaiseOnMessage:=aRaise; + ConA:=A.Client.Connection; + ConB:=B.Client.Connection; + + A.Client.SendMessage('bye'); { the peer closes } + B.Client.SendMessage('echo'); { the peer answers, so a message waits } + Sleep(500); + + { Both have to be waiting before the pump's first pass, otherwise A can be + notified in one pass and B raise in another, and every assertion below + would hold without the arrangement ever existing. } + Result.Armed:=(Pump.ClientCount=2) and A.Client.Active and B.Client.Active + and (ConA.CheckIncoming(50,False)=irWaiting) + and (ConB.CheckIncoming(50,False)=irWaiting); + + Pump.Execute; + { Several passes. A notification that were merely late would arrive here. } + Sleep(1500); + + Result.Errors:=ReadCounter(Sink.FErrors); + Result.Error:=Sink.LastError; + Result.DiscA:=ReadCounter(A.FDisconnects); + Result.TrackedA:=Pump.Tracks(ConA); + Result.TrackedB:=Pump.Tracks(ConB); + Pump.Terminate; + finally + FreeAndNil(A); + FreeAndNil(B); + FreeAndNil(Pump); + FreeAndNil(Sink); + FreeAndNil(SrvA); + FreeAndNil(SrvB); + end; +end; + +Procedure TestExceptionSkipsNotification; +Var + Ctl, Prov : TSkipRound; +begin + BeginScenario('an exception in a later client and a pending notification'); + + Say(' control round: no failing callback'); + Ctl:=RunSkipRound(False); + Say(Format(' control: errors=%d, OnDisconnect A=%d, A tracked=%s', + [Ctl.Errors,Ctl.DiscA,BoolToStr(Ctl.TrackedA,True)])); + Check('control: the arrangement is what it claims to be',Ctl.Armed, + 'the pump did not have two live clients before the pass'); + Check('control: the closed connection is reported',Ctl.DiscA>=1, + Format('OnDisconnect ran %d time(s)',[Ctl.DiscA])); + Check('control: no error is reported',Ctl.Errors=0,Ctl.Error); + + Say(' provocation round: the second client''s callback raises'); + Prov:=RunSkipRound(True); + Say(Format(' provocation: errors=%d, OnDisconnect A=%d, ' + +'A tracked=%s, B tracked=%s', + [Prov.Errors,Prov.DiscA,BoolToStr(Prov.TrackedA,True), + BoolToStr(Prov.TrackedB,True)])); + + Check('provocation: the arrangement is what it claims to be',Prov.Armed, + 'the two clients were not both waiting before the pass, so A and B ' + +'may well have been handled in different passes'); + Check('the failing callback is reported through OnError', + (Prov.Errors>=1) and (Pos('deliberate failure',Prov.Error)>0), + 'the reported error was "'+Prov.Error+'", not the callback''s own'); + Check('the closed connection is still reported to its owner',Prov.DiscA>=1, + Format('OnDisconnect ran %d time(s)',[Prov.DiscA])); + Check('a connection the pump dropped was either reported or is still tracked', + (Prov.DiscA>=1) or Prov.TrackedA, + 'it is in neither state: the owner was not told and no later pass ' + +'can find it again'); + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 13: Terminate called from a synchronized callback. + + An OnMessage callback hands the message to the main thread with + TThread.Synchronize, and the method that runs there stops the pump - the + ordinary shape of "a message arrived, the form decided to close". + + The reader thread is then parked in Synchronize until that method + returns, and the method is inside Terminate waiting for the reader + thread. Servicing CheckSynchronize cannot help: the queue is empty, the + one entry is the method already running. + + The main thread plays the part of an application's message loop and + services the queue itself. + --------------------------------------------------------------------- } +Procedure TestTerminateFromSynchronizedCallback; +Var + Srv : TEchoServer; + Pump : TWSThreadMessagePump; + Cli : TTestClient; + Started, Deadline : QWord; +begin + BeginScenario('Terminate called from a synchronized callback'); + Srv:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TWSThreadMessagePump.Create(Nil); + Cli:=Nil; + try + Srv.Start; + Sleep(150); + Pump.Execute; + Cli:=TTestClient.Create(Srv.Port,Pump); + Cli.SyncTerminatePump:=Pump; + Cli.Client.Connect; + Cli.Client.SendMessage('close'); + + { If Terminate cannot return inside the synchronized method, control + never comes back here and the watchdog reports the hang. } + Started:=TThread.GetTickCount64; + Deadline:=Started+10000; + While (ReadCounter(Cli.FSyncTermReturned)=0) + and (ReadCounter(Cli.FSyncTermRaised)=0) + and (TThread.GetTickCount640, + 'the message never arrived, so this scenario decides nothing'); + Check('Terminate returns when called from a synchronized callback', + ReadCounter(Cli.FSyncTermReturned)>0, + Format('raised=%d',[ReadCounter(Cli.FSyncTermRaised)])); + finally + if Assigned(Cli) then + Cli.SyncTerminatePump:=Nil; + CheckSynchronize(0); + FreeAndNil(Cli); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 14: a control callback that removes an earlier client while + the pump is still iterating. + + ReadConnections walks its list by index and, on irClose, deletes the + entry at the current index. The close frame is dispatched to OnControl + *inside* CheckIncoming, before irClose is returned, and the pump's list + lock is re-entrant on its own thread - so a callback that disconnects + another client removes that client from the very list being walked. + + Three clients in list order A, B, C. Only B has something waiting: a + close frame. B's control callback disconnects A, which sits before B, + so every later entry moves down by one. If the deletion that follows + still uses B's old index, it hits C - a connection nobody closed. + --------------------------------------------------------------------- } +Procedure TestControlCallbackShiftsIndex; +Var + SrvA, SrvB, SrvC : TEchoServer; + Pump : TProbePump; + A, B, C : TTestClient; + ConB, ConC : TWSClientConnection; + Armed, TrackedC, EchoC : Boolean; + MsgC : LongInt; + SendErr : String; +begin + BeginScenario('a close-frame callback removes an earlier client mid-iteration'); + SrvA:=TEchoServer.Create(NextPort,smEcho,False); + SrvB:=TEchoServer.Create(NextPort,smCloseFrameAfterMessage,False); + SrvC:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TProbePump.Create(Nil); + A:=Nil; + B:=Nil; + C:=Nil; + try + SrvA.Start; + SrvB.Start; + SrvC.Start; + Sleep(200); + + { The pump stays stopped until B's close frame is waiting, so that the + first pass is the one that meets it. Connect order is list order. } + A:=TTestClient.Create(SrvA.Port,Pump); + B:=TTestClient.Create(SrvB.Port,Pump); + C:=TTestClient.Create(SrvC.Port,Pump); + A.Client.Connect; + B.Client.Connect; + C.Client.Connect; + ConB:=B.Client.Connection; + ConC:=C.Client.Connection; + B.TearDownOnCloseFrame:=A.Client; + B.Client.SendMessage('bye'); + Sleep(500); + + Armed:=(Pump.ClientCount=3) + and (A.Client.Connection.CheckIncoming(0,False)=irNone) + and (ConB.CheckIncoming(50,False)=irWaiting) + and (ConC.CheckIncoming(0,False)=irNone); + Check('the arrangement is what it claims to be',Armed, + 'three registered clients with only B waiting were not in place'); + if Armed then + begin + Pump.Execute; + WaitForCount(B.FDisconnects,1); + Sleep(500); + + TrackedC:=Pump.Tracks(ConC); + Say(Format(' B close frames=%d, A OnDisconnect=%d, B OnDisconnect=%d', + [ReadCounter(B.FControlCloses),ReadCounter(A.FDisconnects), + ReadCounter(B.FDisconnects)])); + Say(Format(' C: tracked by the pump=%s, Active=%s, OnDisconnect=%d', + [BoolToStr(TrackedC,True),BoolToStr(C.Client.Active,True), + ReadCounter(C.FDisconnects)])); + + if (ReadCounter(B.FControlCloses)=0) or (ReadCounter(A.FDisconnects)=0) then + Check('the provocation actually ran',False, + 'the close-frame callback of B did not remove A, so the list ' + +'never shifted and this round decides nothing') + else + begin + Check('B, whose peer closed, is reported',ReadCounter(B.FDisconnects)>=1); + Check('C, which nobody closed, is still tracked by the pump',TrackedC, + Format('C is gone from the pump but Active=%s with %d disconnect(s)', + [BoolToStr(C.Client.Active,True),ReadCounter(C.FDisconnects)])); + + { The practical consequence: does C still get its messages? } + MsgC:=ReadCounter(C.FMessages); + SendErr:=''; + try + C.Client.SendMessage('still there?'); + except + On E : Exception do + SendErr:=E.ClassName+': '+E.Message; + end; + EchoC:=(SendErr='') and WaitForCount(C.FMessages,MsgC+1,2000); + if SendErr<>'' then + Check('C still receives its echo',False,'send raised '+SendErr) + else + Check('C still receives its echo',EchoC, + 'the echo was sent back but nobody reads C any more'); + end; + end; + finally + FreeAndNil(A); + FreeAndNil(B); + FreeAndNil(C); + FreeAndNil(Pump); + FreeAndNil(SrvA); + FreeAndNil(SrvB); + FreeAndNil(SrvC); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 15: the peer goes away in the middle of a frame. + + The frame header announces five payload bytes, and then the peer closes + the connection instead of sending them - once gracefully, once with a + reset. The payload read therefore fails with an ordinary EReadError, + which is neither EWSReadInterrupted nor irClose. + + Asked: is the owner told, does the connection leave the pump, and does + the pump stop reporting errors for it once that has happened. + --------------------------------------------------------------------- } + +{ Test servers accept in threads of their own and may not accept yet when + a client connects - seen on linux for the stall server and for the echo + server alike. A refused connect is a setup race, not a result, so it is + retried for up to two seconds before it counts. Retries are logged, so a + run shows where they happened. } +Procedure ConnectWithRetry(aClient : TCustomWebsocketClient); +Var + Tries : Integer; +begin + Tries:=0; + repeat + try + aClient.Connect; + Exit; + except + On E : ESocketError do + begin + Inc(Tries); + Say(Format(' connect to port %d refused, retry %d',[aClient.Port,Tries])); + if Tries>=20 then + Raise; + Sleep(100); + end; + end; + until False; +end; + +Type + TMidFrameRound = Record + Armed, Active, Tracked : Boolean; + Disc, Messages, Errors, ReadErrors, LateErrors : LongInt; + NotifiedMs : Int64; + FirstError : String; + end; + +Function RunMidFrameClose(aMode : LongInt) : TMidFrameRound; +Var + Srv : TStallServer; + Pump : TProbePump; + Sink : TErrorSink; + Cli : TTestClient; + Con : TWSClientConnection; + Port : Word; + Started : QWord; + AtNotify : LongInt; +begin + Result.Armed:=False; + Result.Active:=False; + Result.Tracked:=False; + Result.Disc:=0; + Result.Messages:=0; + Result.Errors:=0; + Result.ReadErrors:=0; + Result.LateErrors:=0; + Result.NotifiedMs:=-1; + Result.FirstError:=''; + Port:=NextPort; + Srv:=TStallServer.Create(Port); + Srv.AfterHalf:=aMode; + Pump:=TProbePump.Create(Nil); + Sink:=TErrorSink.Create; + Cli:=Nil; + try + Pump.OnError:=@Sink.DoError; + Sleep(250); + Pump.Execute; + Cli:=TTestClient.Create(Port,Pump,False,True); + ConnectWithRetry(Cli.Client); + Con:=Cli.Client.Connection; + Result.Armed:=WaitForCount(Srv.FHalfSent,1) and WaitForCount(Srv.FClosedAfterHalf,1); + + Started:=TThread.GetTickCount64; + if WaitForCount(Cli.FDisconnects,1,3000) then + Result.NotifiedMs:=TThread.GetTickCount64-Started; + { The pump may report the error just after the disconnect, so let that + one arrive first. A further second then: an error that keeps coming + would mean the pump is still reading the dead connection. } + WaitForCount(Sink.FErrors,1,1000); + AtNotify:=ReadCounter(Sink.FErrors); + Sleep(1000); + Result.LateErrors:=ReadCounter(Sink.FErrors)-AtNotify; + Result.Errors:=ReadCounter(Sink.FErrors); + Result.ReadErrors:=ReadCounter(Sink.FReadErrors); + Result.FirstError:=Sink.LastError; + Result.Disc:=ReadCounter(Cli.FDisconnects); + Result.Messages:=ReadCounter(Cli.FMessages); + Result.Active:=Cli.Client.Active; + Result.Tracked:=Pump.Tracks(Con); + Pump.Terminate; + finally + FreeAndNil(Cli); + FreeAndNil(Pump); + FreeAndNil(Sink); + if Assigned(Srv) then + begin + Srv.Shutdown; + FreeAndNil(Srv); + end; + end; +end; + +Procedure TestPeerLeavesMidFrame; + + Procedure Report(Const aName : String; Const R : TMidFrameRound; + aExpectReadError : Boolean); + begin + Say(Format(' %s: OnDisconnect=%d after %d ms, messages=%d, Active=%s, ' + +'tracked=%s, OnError=%d (EWSReadError=%d, %d in the second ' + +'after), last error: %s', + [aName,R.Disc,R.NotifiedMs,R.Messages,BoolToStr(R.Active,True), + BoolToStr(R.Tracked,True),R.Errors,R.ReadErrors,R.LateErrors, + R.FirstError])); + if not R.Armed then + begin + Check(aName+': the peer left mid-frame',False, + 'the half frame or the close did not happen, so this round decides nothing'); + Exit; + end; + Check(aName+': the owner is told',R.Disc=1, + Format('OnDisconnect ran %d time(s)',[R.Disc])); + Check(aName+': the terminal read is reported before the bounded wait expires', + R.NotifiedMs>=0, + Format('notification time=%d ms',[R.NotifiedMs])); + Check(aName+': an incomplete frame is never delivered',R.Messages=0, + Format('%d message callback(s)',[R.Messages])); + Check(aName+': the client becomes inactive',not R.Active, + 'Active='+BoolToStr(R.Active,True)); + Check(aName+': the connection leaves the pump',not R.Tracked); + if aExpectReadError then + begin + Check(aName+': reset reports exactly one terminal read error', + (R.Errors=1) and (R.ReadErrors=1), + Format('OnError=%d, EWSReadError=%d', + [R.Errors,R.ReadErrors])); + end + else + begin + Check(aName+': orderly EOF does not report a transport error', + (R.Errors=0) and (R.ReadErrors=0), + Format('OnError=%d, EWSReadError=%d', + [R.Errors,R.ReadErrors])); + end; + Check(aName+': no errors keep arriving for it afterwards',R.LateErrors=0, + Format('%d more in one second',[R.LateErrors])); + end; + +Var + R : TMidFrameRound; +begin + BeginScenario('the peer goes away in the middle of a frame'); + R:=RunMidFrameClose(1); + Report('graceful close',R,False); + R:=RunMidFrameClose(2); + Report('reset',R,True); + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 16: an OnError handler that hands the error to the main thread, + which then disconnects another client of the same pump. + + Reporting a failed read through OnError is ordinary, and so is passing + it to the main thread with TThread.Synchronize. If the pump still holds + its list lock while calling OnError, the synchronized method blocks in + RemoveClient on that lock, while the reader waits in Synchronize for the + method to return. + + Client X has its connection reset in the middle of a frame, which produces + the error; client Y is healthy and is the one the main thread disconnects. + --------------------------------------------------------------------- } + +Type + TSyncErrorSink = Class + Private + FErrors : LongInt; + FEntered : LongInt; // the synchronized method started + FReturned : LongInt; // the synchronized method returned + FTarget : TCustomWebsocketClient; + Procedure DisconnectTarget; + Public + Procedure DoError(Sender : TObject; E : Exception); + end; + +Procedure TSyncErrorSink.DisconnectTarget; +Var + T : TCustomWebsocketClient; +begin + BumpCounter(FEntered); + T:=FTarget; + FTarget:=Nil; + if Assigned(T) then + T.Disconnect(False); + BumpCounter(FReturned); +end; + +Procedure TSyncErrorSink.DoError(Sender : TObject; E : Exception); +begin + BumpCounter(FErrors); + { Only the first error is handed over. } + if (ReadCounter(FErrors)=1) and Assigned(FTarget) then + TThread.Synchronize(TThread.CurrentThread,@DisconnectTarget); +end; + +Procedure TestOnErrorSynchronizesDisconnect; +Var + StallSrv : TStallServer; + EchoSrv : TEchoServer; + Pump : TWSThreadMessagePump; + Sink : TSyncErrorSink; + X, Y : TTestClient; + Port : Word; + Started, Deadline : QWord; +begin + BeginScenario('OnError synchronizes a disconnect of another client'); + Port:=NextPort; + StallSrv:=TStallServer.Create(Port); + { An orderly EOF is intentionally not an OnError (scenario 15 verifies + that policy). Use an abortive close so this scenario really exercises + the OnError callback path it claims to test. } + StallSrv.AfterHalf:=2; + EchoSrv:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TWSThreadMessagePump.Create(Nil); + Sink:=TSyncErrorSink.Create; + X:=Nil; + Y:=Nil; + try + Pump.OnError:=@Sink.DoError; + EchoSrv.Start; + Sleep(250); + { Start the pump before either connection is registered. This ensures its + per-session reader is active while X is reset mid-frame instead of + depending on whether a previously closed socket is reported readable. } + Pump.Execute; + X:=TTestClient.Create(Port,Pump); + Y:=TTestClient.Create(EchoSrv.Port,Pump); + { Make the healthy target active and publish it to the error sink before + arming X. The reset can then never beat FTarget initialization. } + ConnectWithRetry(Y.Client); + Sink.FTarget:=Y.Client; + ConnectWithRetry(X.Client); + + { The main thread plays the part of an application's message loop. If + the handler cannot complete, control never returns here and the + watchdog reports the hang. The limit covers the unpatched frame + reader, which retries for about 1.5 s on Windows first. } + Started:=TThread.GetTickCount64; + Deadline:=Started+12000; + While (ReadCounter(Sink.FReturned)=0) and (TThread.GetTickCount640, + StallSrv.LastError); + Check('the error reached OnError',ReadCounter(Sink.FErrors)>0, + 'no error was reported, so this scenario decides nothing'); + Check('the synchronized disconnect completes',ReadCounter(Sink.FReturned)>0, + Format('entered=%d',[ReadCounter(Sink.FEntered)])); + Check('Y is told it was disconnected',ReadCounter(Y.FDisconnects)>=1); + finally + Sink.FTarget:=Nil; + CheckSynchronize(0); + FreeAndNil(X); + FreeAndNil(Y); + FreeAndNil(Pump); + FreeAndNil(Sink); + FreeAndNil(EchoSrv); + if Assigned(StallSrv) then + begin + StallSrv.Shutdown; + FreeAndNil(StallSrv); + end; + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 17: an OnMessage handler that hands the message to the main + thread, which then disconnects another client of the same pump. + + The same shape as scenario 16, but from OnMessage - by far the more + common place for it. Message callbacks run inside CheckIncoming, while + the pump holds its list lock; the synchronized method's Disconnect needs + that lock in RemoveClient. + --------------------------------------------------------------------- } +Procedure TestOnMessageSynchronizesDisconnect; +Var + SrvA, SrvB : TEchoServer; + Pump : TWSThreadMessagePump; + A, B : TTestClient; + Started, Deadline : QWord; +begin + BeginScenario('OnMessage synchronizes a disconnect of another client'); + SrvA:=TEchoServer.Create(NextPort,smEcho,False); + SrvB:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TWSThreadMessagePump.Create(Nil); + A:=Nil; + B:=Nil; + try + SrvA.Start; + SrvB.Start; + Sleep(250); + A:=TTestClient.Create(SrvA.Port,Pump); + B:=TTestClient.Create(SrvB.Port,Pump); + ConnectWithRetry(A.Client); + ConnectWithRetry(B.Client); + A.SyncDisconnectOther:=B.Client; + Pump.Execute; + A.Client.SendMessage('close the other one'); + + { The main thread plays the part of an application's message loop. If + the method cannot complete, control never returns here and the + watchdog reports the hang. } + Started:=TThread.GetTickCount64; + Deadline:=Started+10000; + While (ReadCounter(A.FSyncOtherReturned)=0) and (TThread.GetTickCount640, + 'no echo arrived, so this scenario decides nothing'); + Check('the synchronized disconnect completes',ReadCounter(A.FSyncOtherReturned)>0, + Format('entered=%d',[ReadCounter(A.FSyncOtherEntered)])); + Check('B is told it was disconnected',ReadCounter(B.FDisconnects)>=1); + finally + if Assigned(A) then + A.SyncDisconnectOther:=Nil; + CheckSynchronize(0); + FreeAndNil(A); + FreeAndNil(B); + FreeAndNil(Pump); + FreeAndNil(SrvA); + FreeAndNil(SrvB); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 18: a client freed from the main thread while its own message + callback is still running on the reader thread. + + Freeing the client disconnects it, and RemoveClient returns before the + connection object is freed. That must not happen while the reader is + still inside the connection's callback. While the pump held its list lock + during callbacks this was protected as a side effect; without the lock it + has to be protected on purpose. Expected to pass everywhere. + + The probe connection keeps its memory after destruction, so a callback + that outlives its connection is counted instead of crashing. + --------------------------------------------------------------------- } +Procedure TestFreeClientDuringItsCallback; +Var + Srv : TEchoServer; + Pump : TWSThreadMessagePump; + A : TTestClient; + Started : QWord; + FreedMs : Int64; +begin + BeginScenario('a client freed while its own callback is still running'); + Srv:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TWSThreadMessagePump.Create(Nil); + A:=Nil; + try + Srv.Start; + Sleep(250); + A:=TTestClient.Create(Srv.Port,Pump,False,False,True); + ConnectWithRetry(A.Client); + A.FSlowConnection:=A.Client.Connection; + A.SlowMessageMs:=800; + Pump.Execute; + A.Client.SendMessage('take your time'); + + Check('the callback started',WaitForCount(A.FSlowEntered,1), + 'no echo arrived, so this scenario decides nothing'); + { Free only the websocket client; the test wrapper and its counters stay. } + Started:=TThread.GetTickCount64; + FreeAndNil(A.FClient); + FreedMs:=TThread.GetTickCount64-Started; + WaitForCount(A.FSlowDone,1,3000); + Say(Format(' freeing the client took %d ms; callback done=%d, saw its connection destroyed=%d', + [FreedMs,ReadCounter(A.FSlowDone),ReadCounter(A.FSlowSawDestroyed)])); + + Check('the callback finished',ReadCounter(A.FSlowDone)>0); + Check('the connection outlives its running callback', + ReadCounter(A.FSlowSawDestroyed)=0, + 'the connection was destroyed while its callback was still running'); + finally + Pump.Terminate; + FreeAndNil(A); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 19: a client freed from the main thread while its own message + callback waits in TThread.Synchronize. + + The callback waits for the main thread; the main thread, inside Free, + waits for the reader to be done with that connection. Only servicing + Synchronize while waiting lets both finish. Nothing in the synchronized + method touches the client. + --------------------------------------------------------------------- } +Procedure TestFreeClientWhileItsCallbackSynchronizes; +Var + Srv : TEchoServer; + Pump : TWSThreadMessagePump; + A : TTestClient; + Started : QWord; + FreedMs : Int64; +begin + BeginScenario('a client freed while its callback waits in Synchronize'); + Srv:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TWSThreadMessagePump.Create(Nil); + A:=Nil; + try + Srv.Start; + Sleep(250); + A:=TTestClient.Create(Srv.Port,Pump,False,False,True); + ConnectWithRetry(A.Client); + A.SyncOnMessage:=True; + Pump.Execute; + A.Client.SendMessage('wait for the main thread'); + + Check('the callback reached Synchronize',WaitForCount(A.FSyncEntered,1), + 'no echo arrived, so this scenario decides nothing'); + { Deliberately not servicing Synchronize here: the question is whether + Free does. If it cannot, control never returns and the watchdog + reports the hang. } + Started:=TThread.GetTickCount64; + FreeAndNil(A.FClient); + FreedMs:=TThread.GetTickCount64-Started; + CheckSynchronize(0); + Say(Format(' freeing the client took %d ms; Synchronize returned=%d', + [FreedMs,ReadCounter(A.FSyncReturned)])); + Check('freeing a client whose callback waits for the main thread returns',True); + Check('the callback''s Synchronize completed',WaitForCount(A.FSyncReturned,1,2000)); + finally + CheckSynchronize(0); + Pump.Terminate; + FreeAndNil(A); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 20: the owner frees a client while the pump's OnDisconnect + notification for that client is still running. + + The peer closes; the pump removes the connection and notifies the owner, + whose OnDisconnect takes a while. The client is already inactive then, so + its destructor's Disconnect returns at once - the question is whether the + connection is freed under the running notification. + --------------------------------------------------------------------- } +Procedure TestFreeClientDuringItsDisconnectNotification; +Var + Srv : TEchoServer; + Pump : TWSThreadMessagePump; + A : TTestClient; + Started : QWord; + FreedMs : Int64; +begin + BeginScenario('a client freed while its OnDisconnect notification runs'); + Srv:=TEchoServer.Create(NextPort,smCloseAfterMessage,False); + Pump:=TWSThreadMessagePump.Create(Nil); + A:=Nil; + try + Srv.Start; + Sleep(250); + A:=TTestClient.Create(Srv.Port,Pump,False,False,True); + ConnectWithRetry(A.Client); + A.FSlowConnection:=A.Client.Connection; + A.FSlowDisconnectMs:=800; + Pump.Execute; + A.Client.SendMessage('bye'); + + Check('the disconnect notification started',WaitForCount(A.FSlowDiscEntered,1), + 'the peer close was not reported, so this scenario decides nothing'); + Started:=TThread.GetTickCount64; + FreeAndNil(A.FClient); + FreedMs:=TThread.GetTickCount64-Started; + WaitForCount(A.FSlowDiscDone,1,3000); + Say(Format(' freeing the client took %d ms; notification done=%d, saw its connection destroyed=%d', + [FreedMs,ReadCounter(A.FSlowDiscDone),ReadCounter(A.FSlowDiscSawDestroyed)])); + Check('the notification finished',ReadCounter(A.FSlowDiscDone)>0); + Check('the connection outlives the running notification', + ReadCounter(A.FSlowDiscSawDestroyed)=0, + 'the connection was destroyed while OnDisconnect for it was still running'); + finally + Pump.Terminate; + FreeAndNil(A); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Helpers for scenarios 21-32: waiting while playing the main thread's + message loop, so that synchronized methods can run. + --------------------------------------------------------------------- } +Function WaitServicing(Var aCounter : LongInt; aWanted : LongInt; + aLimitMs : Integer = WaitLimitMs) : Boolean; +Var + Deadline : QWord; +begin + Deadline:=TThread.GetTickCount64+QWord(aLimitMs); + While (ReadCounter(aCounter)=aWanted; +end; + +Function WaitDestroyed(aAddr : Pointer; aLimitMs : Integer = WaitLimitMs) : Boolean; +Var + Deadline : QWord; +begin + Deadline:=TThread.GetTickCount64+QWord(aLimitMs); + While (not WasDestroyed(aAddr)) and (TThread.GetTickCount640, + 'no echo arrived, so this scenario decides nothing'); + Check('the synchronized disconnect returns',ReadCounter(A.FSelfReturned)=1,A.SelfError); + Check('the callback goes on after Synchronize',ReadCounter(A.FCallbackDone)=1); + Check('OnDisconnect exactly once',ReadCounter(A.FDisconnects)=1, + IntToStr(ReadCounter(A.FDisconnects))); + Check('the client is inactive and no longer tracked', + (not A.Client.Active) and (Pump.ClientCount=0)); + Check('the connection is destroyed exactly once',DestroyedTimes(C0)=1, + IntToStr(DestroyedTimes(C0))); + Check('nothing ran on the connection after its destruction',ReadCounter(ContinuedAfterFree)=0); + finally + CheckSynchronize(0); + Pump.Terminate; + FreeAndNil(A); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 22: the same method disconnects and reconnects the client. + + Besides the lifetime of the old connection this asks whether the + replacement is registered and served, and whether anything the old + connection still does changes the new one's state. + --------------------------------------------------------------------- } +Procedure TestSelfReconnectFromSynchronizedMessage; +Var + Srv : TEchoServer; + Pump : TProbePump; + A : TTestClient; + C0, C1 : TWebSocketClientConnection; +begin + BeginScenario('a synchronized method reconnects the client whose callback waits for it'); + Srv:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TProbePump.Create(Nil); + A:=Nil; + try + Srv.Start; + Sleep(250); + A:=TTestClient.Create(Srv.Port,Pump,False,False,True); + ConnectWithRetry(A.Client); + C0:=A.Client.Connection; + InterLockedExchange(A.FSyncSelfAction,2); + Pump.Execute; + A.Client.SendMessage('reconnect me from the main thread'); + + WaitServicing(A.FSelfEntered,1); + WaitServicing(A.FCallbackDone,1); + WaitDestroyed(C0,2000); + C1:=A.Client.Connection; + if ReadCounter(A.FSelfReturned)=1 then + begin + A.Client.SendMessage('over the new connection'); + WaitServicing(A.FMessages,2); + end; + WaitServicing(A.FDisconnects,2,300); + Say(Format(' method returned=%d raised=%d, callback done=%d, OnDisconnect=%d, Active=%s, new tracked=%s, old destroyed %d times, new destroyed=%s, messages=%d, work after free=%d', + [ReadCounter(A.FSelfReturned),ReadCounter(A.FSelfRaised),ReadCounter(A.FCallbackDone), + ReadCounter(A.FDisconnects),BoolToStr(A.Client.Active,True), + BoolToStr(Assigned(C1) and Pump.Tracks(C1),True),DestroyedTimes(C0), + BoolToStr(WasDestroyed(C1),True),ReadCounter(A.FMessages),ReadCounter(ContinuedAfterFree)])); + Check('the message reached OnMessage',ReadCounter(A.FMessages)>0, + 'no echo arrived, so this scenario decides nothing'); + Check('the synchronized reconnect returns',ReadCounter(A.FSelfReturned)=1,A.SelfError); + Check('the callback goes on after Synchronize',ReadCounter(A.FCallbackDone)=1); + Check('OnDisconnect exactly once, for the old connection',ReadCounter(A.FDisconnects)=1, + IntToStr(ReadCounter(A.FDisconnects))); + Check('the client is active on a new, tracked connection', + A.Client.Active and Assigned(C1) and (C1<>C0) and Pump.Tracks(C1)); + Check('the old connection is destroyed exactly once',DestroyedTimes(C0)=1, + IntToStr(DestroyedTimes(C0))); + Check('the new connection is alive',Assigned(C1) and not WasDestroyed(C1)); + Check('the new connection is served',ReadCounter(A.FMessages)>=2); + Check('nothing ran on a connection after its destruction',ReadCounter(ContinuedAfterFree)=0); + finally + CheckSynchronize(0); + Pump.Terminate; + FreeAndNil(A); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenarios 23 and 24: the peer closes, and OnDisconnect reconnects the + client - through a synchronized method (23) or directly on the reader + thread (24). The pump is still inside its notification for the old + connection while the replacement is created and registered. + --------------------------------------------------------------------- } +Procedure RunReconnectFromDisconnect(aMode : LongInt; Const aName : String); +Var + Srv : TEchoServer; + Pump : TProbePump; + A : TTestClient; + C0, C1 : TWebSocketClientConnection; +begin + BeginScenario(aName); + Srv:=TEchoServer.Create(NextPort,smCloseAfterMessage,False); + Pump:=TProbePump.Create(Nil); + A:=Nil; + try + Srv.Start; + Sleep(250); + A:=TTestClient.Create(Srv.Port,Pump,False,False,True); + ConnectWithRetry(A.Client); + C0:=A.Client.Connection; + InterLockedExchange(A.FReconnectOnDisconnect,aMode); + Pump.Execute; + A.Client.SendMessage('close, then reconnect me'); + + WaitServicing(A.FSelfEntered,1); + WaitServicing(A.FSelfReturned,1); + WaitDestroyed(C0,2000); + WaitServicing(A.FDisconnects,2,300); + C1:=A.Client.Connection; + Say(Format(' reconnect entered=%d returned=%d raised=%d, OnDisconnect=%d, Active=%s, new tracked=%s, old destroyed %d times', + [ReadCounter(A.FSelfEntered),ReadCounter(A.FSelfReturned),ReadCounter(A.FSelfRaised), + ReadCounter(A.FDisconnects),BoolToStr(A.Client.Active,True), + BoolToStr(Assigned(C1) and Pump.Tracks(C1),True),DestroyedTimes(C0)])); + Check('the peer close was reported',ReadCounter(A.FDisconnects)>=1, + 'no OnDisconnect, so this scenario decides nothing'); + if ReadCounter(A.FDisconnects)>=1 then + begin + Check('the reconnect returns',ReadCounter(A.FSelfReturned)=1,A.SelfError); + Check('OnDisconnect exactly once, for the old connection',ReadCounter(A.FDisconnects)=1, + IntToStr(ReadCounter(A.FDisconnects))); + Check('the client is active on a new, tracked connection', + A.Client.Active and Assigned(C1) and (C1<>C0) and Pump.Tracks(C1)); + Check('the old connection is destroyed exactly once',DestroyedTimes(C0)=1, + IntToStr(DestroyedTimes(C0))); + Check('nothing ran on a connection after its destruction',ReadCounter(ContinuedAfterFree)=0); + if ReadCounter(A.FSelfReturned)=1 then + begin + A.Client.SendMessage('close again'); + Check('the new connection is served: its close is reported',WaitServicing(A.FDisconnects,2)); + end; + end; + finally + CheckSynchronize(0); + Pump.Terminate; + FreeAndNil(A); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +Procedure TestReconnectFromSynchronizedDisconnect; +begin + RunReconnectFromDisconnect(2,'OnDisconnect synchronizes a reconnect of the same client'); +end; + +Procedure TestReconnectFromDisconnectOnReader; +begin + RunReconnectFromDisconnect(1,'OnDisconnect reconnects the same client on the reader thread'); +end; + +{ --------------------------------------------------------------------- + Scenarios 25 and 26: a client disconnects itself from its own callback + on the reader thread - from OnMessage (25) and from the control callback + for a close frame (26). The pump is still inside CheckIncoming for that + connection: HandleIncoming goes on after the callback, and for a close + frame it sends the close reply. A second client shows whether the pump + keeps serving. + --------------------------------------------------------------------- } +Procedure RunSelfDisconnectOnReader(aCloseFrame : Boolean; Const aName : String); +Var + SrvA, SrvB : TEchoServer; + Pump : TProbePump; + A, B : TTestClient; + C0 : TWebSocketClientConnection; +begin + BeginScenario(aName); + if aCloseFrame then + SrvA:=TEchoServer.Create(NextPort,smCloseFrameAfterMessage,False) + else + SrvA:=TEchoServer.Create(NextPort,smEcho,False); + SrvB:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TProbePump.Create(Nil); + A:=Nil; + B:=Nil; + try + SrvA.Start; + SrvB.Start; + Sleep(250); + A:=TTestClient.Create(SrvA.Port,Pump,False,False,True); + B:=TTestClient.Create(SrvB.Port,Pump); + ConnectWithRetry(A.Client); + ConnectWithRetry(B.Client); + C0:=A.Client.Connection; + if aCloseFrame then + A.FCloseFrameSelfDisconnect:=True + else + InterLockedExchange(A.FDirectSelfAction,1); + Pump.Execute; + A.Client.SendMessage('disconnect me from my own callback'); + + WaitForCount(A.FSelfEntered,1); + WaitDestroyed(C0,2000); + Sleep(300); // room for a second notification + B.Client.SendMessage('still served?'); + WaitForCount(B.FMessages,1); + Say(Format(' action entered=%d returned=%d raised=%d, close frames=%d, OnDisconnect=%d, destroyed %d times, work after free=%d, other client messages=%d', + [ReadCounter(A.FSelfEntered),ReadCounter(A.FSelfReturned),ReadCounter(A.FSelfRaised), + ReadCounter(A.FControlCloses),ReadCounter(A.FDisconnects),DestroyedTimes(C0), + ReadCounter(ContinuedAfterFree),ReadCounter(B.FMessages)])); + Check('the callback ran',ReadCounter(A.FSelfEntered)=1, + 'no echo or close frame arrived, so this scenario decides nothing'); + Check('the self disconnect returns',ReadCounter(A.FSelfReturned)=1,A.SelfError); + if aCloseFrame then + Check('exactly one close control event',ReadCounter(A.FControlCloses)=1, + IntToStr(ReadCounter(A.FControlCloses))); + Check('OnDisconnect exactly once',ReadCounter(A.FDisconnects)=1, + IntToStr(ReadCounter(A.FDisconnects))); + Check('the client is inactive and no longer tracked', + (not A.Client.Active) and not Pump.Tracks(C0)); + Check('the connection is destroyed exactly once',DestroyedTimes(C0)=1, + IntToStr(DestroyedTimes(C0))); + Check('nothing ran on the connection after its destruction',ReadCounter(ContinuedAfterFree)=0, + 'HandleIncoming went on, or a reply was sent, on a destroyed connection'); + Check('the pump still serves the other client',ReadCounter(B.FMessages)>=1); + finally + Pump.Terminate; + FreeAndNil(A); + FreeAndNil(B); + FreeAndNil(Pump); + FreeAndNil(SrvA); + FreeAndNil(SrvB); + end; + EndScenario; +end; + +Procedure TestSelfDisconnectFromOwnMessage; +begin + RunSelfDisconnectOnReader(False,'a client disconnects itself from its own OnMessage'); +end; + +Procedure TestSelfDisconnectFromCloseFrame; +begin + RunSelfDisconnectOnReader(True,'a client disconnects itself from its close-frame callback'); +end; + +{ --------------------------------------------------------------------- + Scenarios 27 and 28: a client whose frame read stalls is freed (27) or + disconnected (28) while the pump keeps running - no Terminate first. + The call must return, and another client of the same pump must still be + served afterwards. + --------------------------------------------------------------------- } +Procedure RunStalledClientWithoutTerminate(aFree : Boolean; Const aName : String); +Var + Srv : TStallServer; + SrvB : TEchoServer; + Pump : TWSThreadMessagePump; + Cli, B : TTestClient; + Started, Elapsed : QWord; + Port : Word; + Waited : Integer; +begin + BeginScenario(aName); + Port:=NextPort; + Srv:=TStallServer.Create(Port); + SrvB:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TWSThreadMessagePump.Create(Nil); + Cli:=Nil; + B:=Nil; + try + SrvB.Start; + Sleep(250); + B:=TTestClient.Create(SrvB.Port,Pump); + ConnectWithRetry(B.Client); + Pump.Execute; + Cli:=TTestClient.Create(Port,Pump,False,True); // instrumented reads + ConnectWithRetry(Cli.Client); + + Check('stall peer completed the handshake', + WaitForCount(Srv.FHandshakes,1),Srv.LastError); + Check('stall peer sent the partial frame', + WaitForCount(Srv.FHalfSent,1),Srv.LastError); + Waited:=0; + While (not ReaderIsInsideRead) and (Waited=1); + finally + Pump.Terminate; + FreeAndNil(Cli); + FreeAndNil(B); + FreeAndNil(Pump); + if Assigned(Srv) then + begin + Srv.Shutdown; + FreeAndNil(Srv); + end; + FreeAndNil(SrvB); + end; + EndScenario; +end; + +Procedure TestFreeStalledClientWithoutTerminate; +begin + RunStalledClientWithoutTerminate(True,'a client with a stalled read is freed without Terminate'); +end; + +Procedure TestDisconnectStalledClientWithoutTerminate; +begin + RunStalledClientWithoutTerminate(False,'a client with a stalled read is disconnected without Terminate'); +end; + +{ --------------------------------------------------------------------- + Scenario 29: while the pump's OnDisconnect notification for a peer close + is still running, the owner disconnects and reconnects the client. + + The late end of the old connection's notification must neither free the + connection under itself, nor report a second disconnect, nor mark the + reconnected client inactive. + --------------------------------------------------------------------- } +Procedure TestReconnectDuringPumpNotification; +Var + Srv : TEchoServer; + Pump : TProbePump; + A : TTestClient; + C0, C1 : TWebSocketClientConnection; +begin + BeginScenario('the owner reconnects while the pump''s OnDisconnect runs'); + Srv:=TEchoServer.Create(NextPort,smCloseAfterMessage,False); + Pump:=TProbePump.Create(Nil); + A:=Nil; + try + Srv.Start; + Sleep(250); + A:=TTestClient.Create(Srv.Port,Pump,False,False,True); + ConnectWithRetry(A.Client); + C0:=A.Client.Connection; + A.FSlowConnection:=C0; + A.FSlowDisconnectMs:=800; + Pump.Execute; + A.Client.SendMessage('bye'); + + Check('the disconnect notification started',WaitForCount(A.FSlowDiscEntered,1), + 'the peer close was not reported, so this scenario decides nothing'); + if ReadCounter(A.FSlowDiscEntered)>0 then + begin + A.Client.Disconnect(False); + ConnectWithRetry(A.Client); + C1:=A.Client.Connection; + WaitForCount(A.FSlowDiscDone,1,3000); + WaitDestroyed(C0,2000); + WaitForCount(A.FDisconnects,2,300); + Say(Format(' notification done=%d, saw its connection destroyed=%d, OnDisconnect=%d, Active=%s, new tracked=%s, old destroyed %d times', + [ReadCounter(A.FSlowDiscDone),ReadCounter(A.FSlowDiscSawDestroyed), + ReadCounter(A.FDisconnects),BoolToStr(A.Client.Active,True), + BoolToStr(Assigned(C1) and Pump.Tracks(C1),True),DestroyedTimes(C0)])); + Check('the notification finished',ReadCounter(A.FSlowDiscDone)>0); + Check('the connection outlives the running notification', + ReadCounter(A.FSlowDiscSawDestroyed)=0); + Check('OnDisconnect exactly once',ReadCounter(A.FDisconnects)=1, + IntToStr(ReadCounter(A.FDisconnects))); + Check('the reconnected client stays active on a tracked connection', + A.Client.Active and Assigned(C1) and (C1<>C0) and Pump.Tracks(C1)); + Check('the old connection is destroyed exactly once',DestroyedTimes(C0)=1, + IntToStr(DestroyedTimes(C0))); + A.FSlowDisconnectMs:=0; + if A.Client.Active then + begin + A.Client.SendMessage('bye again'); + Check('the new connection is served: its close is reported',WaitForCount(A.FDisconnects,2)); + end; + end; + finally + Pump.Terminate; + FreeAndNil(A); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 30: a worker thread frees a client while the client's callback + waits in Synchronize, and the main thread keeps servicing its queue. + Nothing waits for the worker, so Free may wait for the callback - and + has to, before the client goes. + --------------------------------------------------------------------- } +Type + TFreeClientThread = Class(TThread) + Private + FOwner : TTestClient; + Public + FDone : LongInt; + FSyncReturnedWhenFreed : LongInt; + FElapsedMs : Int64; + Constructor Create(aOwner : TTestClient); + Procedure Execute; override; + end; + +Constructor TFreeClientThread.Create(aOwner : TTestClient); +begin + FOwner:=aOwner; + FreeOnTerminate:=False; + Inherited Create(False); +end; + +Procedure TFreeClientThread.Execute; +Var + Started : QWord; +begin + Started:=TThread.GetTickCount64; + FreeAndNil(FOwner.FClient); + FElapsedMs:=TThread.GetTickCount64-Started; + InterLockedExchange(FSyncReturnedWhenFreed,ReadCounter(FOwner.FSyncReturned)); + BumpCounter(FDone); +end; + +{ Sets a flag shortly after a counter was reached: releases a held callback + once another thread has been observed inside the operation that is meant + to overlap it. After ten seconds it releases anyway, without FFired. } +Type + TReleaseWhen = Class(TThread) + Private + FWatch : PLongInt; + FFlag : PLongInt; + Public + FFired : LongInt; + Constructor Create(aWatch, aFlag : PLongInt); + Procedure Execute; override; + end; + +Constructor TReleaseWhen.Create(aWatch, aFlag : PLongInt); +begin + FWatch:=aWatch; + FFlag:=aFlag; + FreeOnTerminate:=False; + Inherited Create(False); +end; + +Procedure TReleaseWhen.Execute; +Var + Deadline : QWord; +begin + Deadline:=TThread.GetTickCount64+10000; + While (InterLockedExchangeAdd(FWatch^,0)=0) and (TThread.GetTickCount640 then + BumpCounter(FFired); + Sleep(200); + InterLockedExchange(FFlag^,1); +end; + +Procedure TestFreeClientFromWorkerWhileCallbackSynchronizes; +Var + Srv : TEchoServer; + Pump : TWSThreadMessagePump; + A : TTestClient; + W : TFreeClientThread; +begin + BeginScenario('a worker frees a client while its callback waits in Synchronize'); + Srv:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TWSThreadMessagePump.Create(Nil); + A:=Nil; + W:=Nil; + try + Srv.Start; + Sleep(250); + A:=TTestClient.Create(Srv.Port,Pump,False,False,True); + ConnectWithRetry(A.Client); + A.SyncOnMessage:=True; + Pump.Execute; + A.Client.SendMessage('wait for the main thread'); + + Check('the callback reached Synchronize',WaitForCount(A.FSyncEntered,1), + 'no echo arrived, so this scenario decides nothing'); + W:=TFreeClientThread.Create(A); + { The queue is not serviced until the worker is inside Free; otherwise + the callback could finish first and nothing would overlap. } + WaitForCount(ProbeClientDestroying,1,2000); + Sleep(200); + Check('the worker is inside Free while the callback still waits', + (ReadCounter(ProbeClientDestroying)=1) and (ReadCounter(A.FSyncReturned)=0), + Format('destroying=%d, Synchronize returned=%d', + [ReadCounter(ProbeClientDestroying),ReadCounter(A.FSyncReturned)])); + WaitServicing(W.FDone,1); + Say(Format(' worker done=%d after %d ms; Synchronize had returned when Free did=%d', + [ReadCounter(W.FDone),W.FElapsedMs,ReadCounter(W.FSyncReturnedWhenFreed)])); + Check('the worker''s Free returns',ReadCounter(W.FDone)=1); + Check('Free waited for the callback',ReadCounter(W.FSyncReturnedWhenFreed)=1); + finally + CheckSynchronize(0); + if Assigned(W) then + begin + W.WaitFor; + FreeAndNil(W); + end; + Pump.Terminate; + FreeAndNil(A); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 31: while a slow callback of a client's connection runs, the + owner disconnects and reconnects the client and then frees it. + + The running callback belongs to a connection the client no longer holds. + Freeing the client must still wait for it, because the callback runs on + behalf of the client component. + --------------------------------------------------------------------- } +Procedure TestFreeClientAfterReconnectWhileOldCallbackRuns; +Var + Srv : TEchoServer; + Pump : TWSThreadMessagePump; + A : TTestClient; + Started : QWord; + FreedMs : Int64; + R : TReleaseWhen; +begin + BeginScenario('a reconnected client is freed while its old connection''s callback runs'); + R:=Nil; + Srv:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TWSThreadMessagePump.Create(Nil); + A:=Nil; + try + Srv.Start; + Sleep(250); + A:=TTestClient.Create(Srv.Port,Pump,False,False,True); + ConnectWithRetry(A.Client); + A.FSlowConnection:=A.Client.Connection; + A.FSlowComponent:=A.Client; + A.SlowMessageMs:=1; + A.FHoldUntilRelease:=True; + Pump.Execute; + A.Client.SendMessage('take your time'); + + Check('the callback started',WaitForCount(A.FSlowEntered,1), + 'no echo arrived, so this scenario decides nothing'); + A.Client.Disconnect(False); + ConnectWithRetry(A.Client); + Check('the old callback still runs when Free starts',ReadCounter(A.FSlowDone)=0, + 'disconnect or reconnect waited for the callback, so nothing overlaps'); + { The callback is released only once this thread is inside Free. } + R:=TReleaseWhen.Create(@ProbeClientDestroying,@A.FReleaseHold); + Started:=TThread.GetTickCount64; + FreeAndNil(A.FClient); + FreedMs:=TThread.GetTickCount64-Started; + WaitForCount(A.FSlowDone,1,5000); + Say(Format(' freeing the client took %d ms; callback done=%d, saw the client destroyed=%d, saw its connection destroyed=%d', + [FreedMs,ReadCounter(A.FSlowDone),ReadCounter(A.FSlowSawComponentDestroyed), + ReadCounter(A.FSlowSawDestroyed)])); + Check('the callback finished',ReadCounter(A.FSlowDone)>0); + Check('the callback was released only after Free had begun',ReadCounter(R.FFired)=1); + Check('the callback ended by the release, not by its time limit',ReadCounter(A.FHoldTimedOut)=0); + Check('the client outlives the callback of its old connection', + ReadCounter(A.FSlowSawComponentDestroyed)=0); + Check('the old connection outlives its callback',ReadCounter(A.FSlowSawDestroyed)=0); + finally + if Assigned(R) then + begin + R.WaitFor; + FreeAndNil(R); + end; + Pump.Terminate; + FreeAndNil(A); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 32: the pump is freed while a connection its owner already + disconnected is still inside a slow callback on the reader thread. + The connection must be destroyed exactly once, after its callback, and + the client must be freeable afterwards. + --------------------------------------------------------------------- } +Procedure TestPumpFreedWhileReleasedConnectionPending; +Var + Srv : TEchoServer; + Pump : TProbePump; + A : TTestClient; + C0 : TWebSocketClientConnection; + R : TReleaseWhen; +begin + BeginScenario('the pump is freed while a disconnected connection''s callback runs'); + R:=Nil; + Srv:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TProbePump.Create(Nil); + A:=Nil; + try + Srv.Start; + Sleep(250); + A:=TTestClient.Create(Srv.Port,Pump,False,False,True); + ConnectWithRetry(A.Client); + C0:=A.Client.Connection; + A.FSlowConnection:=C0; + A.SlowMessageMs:=1; + A.FHoldUntilRelease:=True; + Pump.Execute; + A.Client.SendMessage('take your time'); + + Check('the callback started',WaitForCount(A.FSlowEntered,1), + 'no echo arrived, so this scenario decides nothing'); + A.Client.Disconnect(False); + Check('the disconnected connection''s callback still runs when the pump is freed', + ReadCounter(A.FSlowDone)=0, + 'Disconnect waited for the callback, so nothing was pending'); + { The callback is released only once this thread is inside the pump's Free. } + R:=TReleaseWhen.Create(@ProbePumpDestroying,@A.FReleaseHold); + FreeAndNil(Pump); + WaitForCount(A.FSlowDone,1,4000); + WaitDestroyed(C0,2000); + { Whether the client still refers to the freed pump is the problem + --pump-first shows; the client is inactive here and does not use it. } + Say(Format(' callback done=%d, saw its connection destroyed=%d, destroyed %d times, pump reference cleared=%s, work after free=%d', + [ReadCounter(A.FSlowDone),ReadCounter(A.FSlowSawDestroyed),DestroyedTimes(C0), + BoolToStr(A.Client.MessagePump=Nil,True),ReadCounter(ContinuedAfterFree)])); + Check('the callback finished',ReadCounter(A.FSlowDone)>0); + Check('the callback was released only after the pump''s Free had begun',ReadCounter(R.FFired)=1); + Check('the callback ended by the release, not by its time limit',ReadCounter(A.FHoldTimedOut)=0); + Check('the connection outlives its callback',ReadCounter(A.FSlowSawDestroyed)=0); + Check('the connection is destroyed exactly once',DestroyedTimes(C0)=1, + IntToStr(DestroyedTimes(C0))); + Check('nothing ran on the connection after its destruction',ReadCounter(ContinuedAfterFree)=0); + FreeAndNil(A); + Check('the client is freed afterwards',True); + finally + if Assigned(R) then + begin + R.WaitFor; + FreeAndNil(R); + end; + if Assigned(Pump) then + Pump.Terminate; + FreeAndNil(A); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 33: a client disconnects itself from its own OnMessage, and its + OnDisconnect handler reconnects it at once - all on the reader thread, + inside the callback of the old connection. + + The reconnect replaces the handshake response. The old connection still + refers to its own response while the callback runs, so that response + must not be freed before the callback returns. + --------------------------------------------------------------------- } +Procedure TestReconnectInsideOwnDisconnect; +Var + Srv : TEchoServer; + Pump : TProbePump; + A : TTestClient; + C0, C1 : TWebSocketClientConnection; +begin + BeginScenario('OnDisconnect reconnects inside the callback that disconnected the client'); + Srv:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TProbePump.Create(Nil); + A:=Nil; + try + Srv.Start; + Sleep(250); + A:=TTestClient.Create(Srv.Port,Pump,False,False,True); + ConnectWithRetry(A.Client); + C0:=A.Client.Connection; + A.FOldResponse:=C0.HandShakeResponse; + InterLockedExchange(A.FReconnectOnDisconnect,1); + InterLockedExchange(A.FDirectSelfAction,1); + Pump.Execute; + A.Client.SendMessage('disconnect me, then reconnect me'); + + WaitForCount(A.FCallbackDone,1); + WaitDestroyed(C0,2000); + C1:=A.Client.Connection; + if Assigned(C1) and A.Client.Active then + begin + A.Client.SendMessage('over the new connection'); + WaitForCount(A.FMessages,2); + end; + Say(Format(' actions entered=%d returned=%d raised=%d, OnDisconnect=%d, Active=%s, new tracked=%s, old response freed during the callback=%d, old destroyed %d times, messages=%d, work after free=%d', + [ReadCounter(A.FSelfEntered),ReadCounter(A.FSelfReturned),ReadCounter(A.FSelfRaised), + ReadCounter(A.FDisconnects),BoolToStr(A.Client.Active,True), + BoolToStr(Assigned(C1) and Pump.Tracks(C1),True),ReadCounter(A.FOldResponseFreed), + DestroyedTimes(C0),ReadCounter(A.FMessages),ReadCounter(ContinuedAfterFree)])); + Check('the message reached OnMessage',ReadCounter(A.FMessages)>0, + 'no echo arrived, so this scenario decides nothing'); + Check('disconnect and reconnect both return',ReadCounter(A.FSelfReturned)=2,A.SelfError); + Check('OnDisconnect exactly once',ReadCounter(A.FDisconnects)=1, + IntToStr(ReadCounter(A.FDisconnects))); + Check('the client is active on a new, tracked connection', + A.Client.Active and Assigned(C1) and (C1<>C0) and Pump.Tracks(C1)); + Check('the old handshake response outlives the callback',ReadCounter(A.FOldResponseFreed)=0); + WaitDestroyed(Pointer(A.FOldResponse),2000); + Check('the old handshake response is destroyed exactly once, afterwards', + DestroyedTimes(Pointer(A.FOldResponse))=1, + IntToStr(DestroyedTimes(Pointer(A.FOldResponse)))); + Check('the old connection is destroyed exactly once',DestroyedTimes(C0)=1, + IntToStr(DestroyedTimes(C0))); + Check('the new connection is served',ReadCounter(A.FMessages)>=2); + Check('nothing ran on a connection after its destruction',ReadCounter(ContinuedAfterFree)=0); + finally + Pump.Terminate; + FreeAndNil(A); + FreeAndNil(Pump); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 34: while a slow callback of a client runs on its pump, the + owner disconnects the client, assigns it another pump and frees it. + + The running callback belongs to the old pump. Either the assignment is + refused while that callback runs, or freeing the client still waits for + it - the client must not be freed under its own callback. + --------------------------------------------------------------------- } +Procedure TestPumpReassignedWhileOldCallbackRuns; +Var + Srv : TEchoServer; + P0, P1 : TWSThreadMessagePump; + A : TTestClient; + Raised, HeldAtAssign : Boolean; + Msg : String; + R : TReleaseWhen; +begin + BeginScenario('the pump is reassigned while a callback runs on the old pump'); + R:=Nil; + Srv:=TEchoServer.Create(NextPort,smEcho,False); + P0:=TWSThreadMessagePump.Create(Nil); + P1:=TWSThreadMessagePump.Create(Nil); + A:=Nil; + try + Srv.Start; + Sleep(250); + A:=TTestClient.Create(Srv.Port,P0,False,False,True); + ConnectWithRetry(A.Client); + A.FSlowConnection:=A.Client.Connection; + A.FSlowComponent:=A.Client; + A.SlowMessageMs:=1; + A.FHoldUntilRelease:=True; + P0.Execute; + A.Client.SendMessage('take your time'); + + Check('the callback started',WaitForCount(A.FSlowEntered,1), + 'no echo arrived, so this scenario decides nothing'); + A.Client.Disconnect(False); + HeldAtAssign:=ReadCounter(A.FSlowDone)=0; + Raised:=False; + Msg:=''; + try + A.Client.MessagePump:=P1; + except + On E : Exception do + begin + Raised:=True; + Msg:=E.ClassName+': '+E.Message; + end; + end; + Check('the callback still runs when the pump is reassigned',HeldAtAssign, + 'Disconnect waited for the callback, so nothing overlaps'); + Check('the reassignment is refused while the old pump runs the callback',Raised); + Check('the refusal is an EWebSocketClient and the old pump stays assigned', + (Pos('EWebSocketClient',Msg)=1) and (A.Client.MessagePump=P0),Msg); + { Released only once this thread is inside Free. } + R:=TReleaseWhen.Create(@ProbeClientDestroying,@A.FReleaseHold); + FreeAndNil(A.FClient); + WaitForCount(A.FSlowDone,1,4000); + Say(Format(' reassignment raised=%s %s; callback done=%d, saw the client destroyed=%d, saw its connection destroyed=%d', + [BoolToStr(Raised,True),Msg,ReadCounter(A.FSlowDone), + ReadCounter(A.FSlowSawComponentDestroyed),ReadCounter(A.FSlowSawDestroyed)])); + Check('the callback finished',ReadCounter(A.FSlowDone)>0); + Check('the callback was released only after Free had begun',ReadCounter(R.FFired)=1); + Check('the callback ended by the release, not by its time limit',ReadCounter(A.FHoldTimedOut)=0); + Check('the client outlives the callback on its old pump', + ReadCounter(A.FSlowSawComponentDestroyed)=0); + Check('the connection outlives its callback',ReadCounter(A.FSlowSawDestroyed)=0); + finally + if Assigned(R) then + begin + R.WaitFor; + FreeAndNil(R); + end; + P0.Terminate; + P1.Terminate; + FreeAndNil(A); + FreeAndNil(P0); + FreeAndNil(P1); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 35: the destructor of a connection its owner released raises + when the reader frees it after the callback. The reader has to go on + serving the other clients of the pump. + --------------------------------------------------------------------- } +Procedure TestConnectionDestructorRaisesInReader; +Var + SrvA, SrvB : TEchoServer; + Pump : TProbePump; + Sink : TErrorSink; + A, B : TTestClient; + C0 : TWebSocketClientConnection; + DiscRaised : Integer; +begin + BeginScenario('a released connection''s destructor raises on the reader thread'); + SrvA:=TEchoServer.Create(NextPort,smEcho,False); + SrvB:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TProbePump.Create(Nil); + Sink:=TErrorSink.Create; + Pump.OnError:=@Sink.DoError; + A:=Nil; + B:=Nil; + try + SrvA.Start; + SrvB.Start; + Sleep(250); + A:=TTestClient.Create(SrvA.Port,Pump,False,False,True); + B:=TTestClient.Create(SrvB.Port,Pump); + ConnectWithRetry(A.Client); + ConnectWithRetry(B.Client); + C0:=A.Client.Connection; + A.SlowMessageMs:=1000; + Pump.Execute; + A.Client.SendMessage('take your time'); + + Check('the callback started',WaitForCount(A.FSlowEntered,1), + 'no echo arrived, so this scenario decides nothing'); + WatchSlowDone:=@A.FSlowDone; + RaiseOnDestroyAddr:=C0; + DiscRaised:=0; + try + A.Client.Disconnect(False); + except + On E : Exception do + Inc(DiscRaised); + end; + WaitForCount(A.FSlowDone,1,3000); + WaitForCount(DestructorRaised,1,3000); + B.Client.SendMessage('still served?'); + WaitForCount(B.FMessages,1); + Say(Format(' destructor raised=%d (raised to the disconnecting caller: %d), pump errors=%d %s, other client messages=%d', + [ReadCounter(DestructorRaised),DiscRaised,Sink.Errors,Sink.LastError,ReadCounter(B.FMessages)])); + Check('the destructor ran and raised',ReadCounter(DestructorRaised)=1, + 'the connection was not freed, so this scenario decides nothing'); + Check('the pump still serves the other client',ReadCounter(B.FMessages)>=1); + Check('the failure does not reach the disconnecting caller',DiscRaised=0); + Check('the destructor ran on the reader thread, after the callback', + (ReadCounter(DestructorRaised)=1) and (DestructorThread<>MainThreadID) + and (CallbackDoneAtDestroy=1), + Format('callback done when the destructor ran=%d',[CallbackDoneAtDestroy])); + Check('the failure is reported through OnError',Sink.Errors>=1); + finally + RaiseOnDestroyAddr:=Nil; + WatchSlowDone:=Nil; + Pump.Terminate; + FreeAndNil(A); + FreeAndNil(B); + FreeAndNil(Pump); + FreeAndNil(Sink); + FreeAndNil(SrvA); + FreeAndNil(SrvB); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 36: OnDisconnect raises while an active client is freed. The + client and its connection must still be destroyed. + --------------------------------------------------------------------- } +Procedure TestOnDisconnectRaisesDuringFree; +Var + Srv : TEchoServer; + Pump : TProbePump; + Sink : TErrorSink; + A : TTestClient; + C0 : TWebSocketClientConnection; + Comp : Pointer; + FreeRaised : Integer; + Msg : String; +begin + BeginScenario('OnDisconnect raises while the client is freed'); + Srv:=TEchoServer.Create(NextPort,smEcho,False); + Pump:=TProbePump.Create(Nil); + Sink:=TErrorSink.Create; + Pump.OnError:=@Sink.DoError; + A:=Nil; + try + Srv.Start; + Sleep(250); + A:=TTestClient.Create(Srv.Port,Pump,False,False,True); + ConnectWithRetry(A.Client); + C0:=A.Client.Connection; + Comp:=Pointer(A.Client); + Pump.Execute; + InterLockedExchange(A.FRaiseOnDisconnect,1); + FreeRaised:=0; + Msg:=''; + try + FreeAndNil(A.FClient); + except + On E : Exception do + begin + Inc(FreeRaised); + Msg:=E.ClassName+': '+E.Message; + end; + end; + WaitDestroyed(C0,2000); + Say(Format(' Free raised=%d %s; OnDisconnect=%d, client destroyed=%s, connection destroyed %d times, pump errors=%d', + [FreeRaised,Msg,ReadCounter(A.FDisconnects),BoolToStr(WasDestroyed(Comp),True), + DestroyedTimes(C0),Sink.Errors])); + Check('OnDisconnect ran',ReadCounter(A.FDisconnects)=1, + 'the handler did not run, so this scenario decides nothing'); + Check('the client is destroyed despite the failing handler',WasDestroyed(Comp)); + Check('its connection is destroyed exactly once',DestroyedTimes(C0)=1, + IntToStr(DestroyedTimes(C0))); + Check('Free does not raise the handler''s exception',FreeRaised=0,Msg); + Check('the failure is reported through OnError',Sink.Errors>=1); + finally + Pump.Terminate; + FreeAndNil(A); + FreeAndNil(Pump); + FreeAndNil(Sink); + FreeAndNil(Srv); + end; + EndScenario; +end; + +{ --------------------------------------------------------------------- + Scenario 37: an incomplete frame on one connection must not monopolize + a shared message pump. + + The stall peer sends a complete frame header but withholds all five bytes + of its announced payload. Read instrumentation establishes that the client + has started processing that frame. Only then is a message sent through a + second, already-proven connection on the same pump. The payload is released + after the result has been captured, so it cannot accidentally unblock a + single-reader implementation early and turn starvation into a pass. + + This deliberately does not require ReaderIsInsideRead: a valid redesign may + use either an independent reader per session or a nonblocking incremental + parser, and the externally observable fairness requirement is the same. + --------------------------------------------------------------------- } +Procedure TestPartialFrameDoesNotStarveSibling; +Var + EchoSrv : TEchoServer; + StallSrv : TStallServer; + Pump : TWSThreadMessagePump; + Healthy, Stalled : TTestClient; + StallPort : Word; + MessageBefore : LongInt; + Started, Elapsed : QWord; + FrameReadStarted, Served : Boolean; + SendErr : String; +begin + BeginScenario('an incomplete frame does not starve a healthy client'); + EchoSrv:=TEchoServer.Create(NextPort,smEcho,False); + StallPort:=NextPort; + StallSrv:=TStallServer.Create(StallPort); + Pump:=TWSThreadMessagePump.Create(Nil); + Healthy:=Nil; + Stalled:=Nil; + try + EchoSrv.Start; + Sleep(250); + Pump.Execute; + + Healthy:=TTestClient.Create(EchoSrv.Port,Pump); + ConnectWithRetry(Healthy.Client); + Healthy.Client.SendMessage('before-stall'); + Check('healthy client works before the partial frame', + WaitForCount(Healthy.FMessages,1) + and (Healthy.LastMessage='before-stall')); + + Stalled:=TTestClient.Create(StallPort,Pump,False,True); + ConnectWithRetry(Stalled.Client); + Check('stall peer completed the handshake', + WaitForCount(StallSrv.FHandshakes,1),StallSrv.LastError); + Check('stall peer sent only the frame header', + WaitForCount(StallSrv.FHalfSent,1),StallSrv.LastError); + + { The instrumented connection is the stalled connection only. Seeing a + read entry therefore proves that the pump has begun consuming its + deliberately incomplete frame, rather than merely having registered + the connection. Give a blocking implementation time to enter its next + payload read; an incremental implementation simply remains idle. } + FrameReadStarted:=WaitForCount(ReadEntries,1,2000); + if FrameReadStarted then + Sleep(100); + Check('the client began processing the incomplete frame',FrameReadStarted, + Format('read entries=%d exits=%d', + [ReadCounter(ReadEntries),ReadCounter(ReadExits)])); + Check('the missing payload is still withheld', + ReadCounter(StallSrv.FRestSent)=0); + + MessageBefore:=ReadCounter(Healthy.FMessages); + SendErr:=''; + Started:=TThread.GetTickCount64; + try + Healthy.Client.SendMessage('during-stall'); + except + On E : Exception do + SendErr:=E.ClassName+': '+E.Message; + end; + Served:=(SendErr='') + and WaitForCount(Healthy.FMessages,MessageBefore+1,2000) + and (Healthy.LastMessage='during-stall'); + Elapsed:=TThread.GetTickCount64-Started; + Say(Format(' sibling echo served=%s after %d ms; stalled reads=%d/%d', + [BoolToStr(Served,True),Elapsed,ReadCounter(ReadEntries), + ReadCounter(ReadExits)])); + + if SendErr<>'' then + Check('the healthy sibling accepts a send while the frame is incomplete', + False,SendErr) + else + Check('the healthy sibling is served while the frame is incomplete', + Served,Format('messages before=%d after=%d, last="%s"', + [MessageBefore,ReadCounter(Healthy.FMessages), + Healthy.LastMessage])); + Check('the stalled payload remained withheld until after the sibling result', + ReadCounter(StallSrv.FRestSent)=0); + Check('the healthy sibling remains active and was not disconnected', + Healthy.Client.Active and (ReadCounter(Healthy.FDisconnects)=0), + Format('Active=%s, OnDisconnect=%d', + [BoolToStr(Healthy.Client.Active,True), + ReadCounter(Healthy.FDisconnects)])); + + { Complete the staged frame only after all fairness observations have + been captured. This keeps cleanup independent of the behavior under + test, even on an implementation that still has one blocking reader. } + StallSrv.SendRest; + WaitForCount(StallSrv.FRestSent,1,2000); + WaitForCount(Stalled.FMessages,1,2000); + Pump.Terminate; + finally + if Assigned(StallSrv) then + StallSrv.SendRest; + FreeAndNil(Healthy); + FreeAndNil(Stalled); + FreeAndNil(Pump); + if Assigned(StallSrv) then + begin + StallSrv.Shutdown; + FreeAndNil(StallSrv); + end; + FreeAndNil(EchoSrv); + end; + EndScenario; +end; + +Const + ScenarioCount = 37; + +Var + ScenarioProcs : Array[1..ScenarioCount] of TScenarioProc; + ScenarioNames : Array[1..ScenarioCount] of String; + +Procedure BuildScenarioTable; +begin + ScenarioProcs[1]:=@TestUpgradeAndEcho; + ScenarioNames[1]:='upgrade handshake and echo'; + ScenarioProcs[2]:=@TestRepeatedExecuteTerminate; + ScenarioNames[2]:='repeated Execute/Terminate'; + ScenarioProcs[3]:=@TestPeerClose; + ScenarioNames[3]:='peer close'; + ScenarioProcs[4]:=@TestTerminateWhileIdle; + ScenarioNames[4]:='Terminate while idle'; + ScenarioProcs[5]:=@TestTLSEchoAndClose; + ScenarioNames[5]:='TLS echo and peer close'; + ScenarioProcs[6]:=@TestPartialFrameStall; + ScenarioNames[6]:='partial frame stall, plain TCP'; + ScenarioProcs[7]:=@TestPartialFrameStallTLS; + ScenarioNames[7]:='partial frame stall, TLS'; + ScenarioProcs[8]:=@TestTerminateFromDisconnect; + ScenarioNames[8]:='Terminate from OnDisconnect'; + ScenarioProcs[9]:=@TestCollateralInterrupt; + ScenarioNames[9]:='collateral interrupt of a healthy client'; + ScenarioProcs[10]:=@TestSynchronizeDuringTerminate; + ScenarioNames[10]:='Terminate while a callback is in Synchronize'; + ScenarioProcs[11]:=@TestFreeQueuedDisconnect; + ScenarioNames[11]:='queued connection destroyed by an earlier callback'; + ScenarioProcs[12]:=@TestExceptionSkipsNotification; + ScenarioNames[12]:='exception in a later client skips a notification'; + ScenarioProcs[13]:=@TestTerminateFromSynchronizedCallback; + ScenarioNames[13]:='Terminate from a synchronized callback'; + ScenarioProcs[14]:=@TestControlCallbackShiftsIndex; + ScenarioNames[14]:='close-frame callback removes an earlier client'; + ScenarioProcs[15]:=@TestPeerLeavesMidFrame; + ScenarioNames[15]:='peer closes or resets in the middle of a frame'; + ScenarioProcs[16]:=@TestOnErrorSynchronizesDisconnect; + ScenarioNames[16]:='OnError synchronizes a disconnect of another client'; + ScenarioProcs[17]:=@TestOnMessageSynchronizesDisconnect; + ScenarioNames[17]:='OnMessage synchronizes a disconnect of another client'; + ScenarioProcs[18]:=@TestFreeClientDuringItsCallback; + ScenarioNames[18]:='client freed while its own callback runs'; + ScenarioProcs[19]:=@TestFreeClientWhileItsCallbackSynchronizes; + ScenarioNames[19]:='client freed while its callback waits in Synchronize'; + ScenarioProcs[20]:=@TestFreeClientDuringItsDisconnectNotification; + ScenarioNames[20]:='client freed while its OnDisconnect notification runs'; + ScenarioProcs[21]:=@TestSelfDisconnectFromSynchronizedMessage; + ScenarioNames[21]:='synchronized method disconnects the waiting client'; + ScenarioProcs[22]:=@TestSelfReconnectFromSynchronizedMessage; + ScenarioNames[22]:='synchronized method reconnects the waiting client'; + ScenarioProcs[23]:=@TestReconnectFromSynchronizedDisconnect; + ScenarioNames[23]:='OnDisconnect synchronizes a reconnect'; + ScenarioProcs[24]:=@TestReconnectFromDisconnectOnReader; + ScenarioNames[24]:='OnDisconnect reconnects on the reader thread'; + ScenarioProcs[25]:=@TestSelfDisconnectFromOwnMessage; + ScenarioNames[25]:='client disconnects itself from its OnMessage'; + ScenarioProcs[26]:=@TestSelfDisconnectFromCloseFrame; + ScenarioNames[26]:='client disconnects itself from its close-frame callback'; + ScenarioProcs[27]:=@TestFreeStalledClientWithoutTerminate; + ScenarioNames[27]:='stalled client freed without Terminate'; + ScenarioProcs[28]:=@TestDisconnectStalledClientWithoutTerminate; + ScenarioNames[28]:='stalled client disconnected without Terminate'; + ScenarioProcs[29]:=@TestReconnectDuringPumpNotification; + ScenarioNames[29]:='owner reconnects while the pump notification runs'; + ScenarioProcs[30]:=@TestFreeClientFromWorkerWhileCallbackSynchronizes; + ScenarioNames[30]:='worker frees a client whose callback synchronizes'; + ScenarioProcs[31]:=@TestFreeClientAfterReconnectWhileOldCallbackRuns; + ScenarioNames[31]:='reconnected client freed while the old callback runs'; + ScenarioProcs[32]:=@TestPumpFreedWhileReleasedConnectionPending; + ScenarioNames[32]:='pump freed while a released connection is pending'; + ScenarioProcs[33]:=@TestReconnectInsideOwnDisconnect; + ScenarioNames[33]:='OnDisconnect reconnects inside the disconnecting callback'; + ScenarioProcs[34]:=@TestPumpReassignedWhileOldCallbackRuns; + ScenarioNames[34]:='pump reassigned while a callback runs on the old pump'; + ScenarioProcs[35]:=@TestConnectionDestructorRaisesInReader; + ScenarioNames[35]:='released connection destructor raises on the reader'; + ScenarioProcs[36]:=@TestOnDisconnectRaisesDuringFree; + ScenarioNames[36]:='OnDisconnect raises while the client is freed'; + ScenarioProcs[37]:=@TestPartialFrameDoesNotStarveSibling; + ScenarioNames[37]:='incomplete frame does not starve a healthy client'; +end; + +{ Run one scenario in this process and exit with its verdict. } +Procedure RunAsChild(aIndex : Integer); +Var + Dog : TWatchdog; +begin + if (aIndex<1) or (aIndex>ScenarioCount) then + begin + Say('no such scenario: '+IntToStr(aIndex)); + Halt(97); + end; + Dog:=TWatchdog.Create(False); + try + ScenarioProcs[aIndex](); + finally + SetDeadline('',0); + Dog.Terminate; + Dog.WaitFor; + Dog.Free; + end; + if ScenariosSkipped>0 then + Halt(2) + else if ScenariosFailed>0 then + Halt(1) + else + Halt(0); +end; + +Type + TChildVerdict = (cvPassed, cvFailed, cvSkipped, cvHung, cvCrashed); + +{ Start ourselves for one scenario and supervise it. The child inherits + our console, so its output appears inline. } +Function RunChild(aIndex : Integer; Out aExit : Integer) : TChildVerdict; +Var + P : TProcess; + Waited : Integer; + Raw : Integer; +begin + aExit:=-1; + Raw:=0; + P:=TProcess.Create(Nil); + try + P.Executable:=ParamStr(0); + P.Parameters.Add('--scenario'); + P.Parameters.Add(IntToStr(aIndex)); + P.Options:=[]; // no pipes: the child writes straight to our console + P.ShowWindow:=swoShow; + P.Execute; + + Waited:=0; + While P.Running and (Waited0) then + begin + aExit:=Raw; + Exit(cvCrashed); // killed by a signal + end; + {$ENDIF} + Case aExit of + 0 : Result:=cvPassed; + 1 : Result:=cvFailed; + 2 : Result:=cvSkipped; + 99: Result:=cvHung; + else + Result:=cvCrashed; + end; + finally + P.Free; + end; +end; + +Var + I, Ex, Bad : Integer; + V : TChildVerdict; + Passed, Failed, Skipped, Hung, Crashed : Integer; + +begin + {$IFDEF UNIX} + { Writing to a socket whose peer has gone away raises SIGPIPE, which kills + the process by default. Every scenario here deliberately shuts sockets + down under a live peer, so ignore it and let write() report EPIPE. } + fpSignal(SIGPIPE,SignalHandler(SIG_IGN)); + {$ENDIF} + InitCriticalSection(OutLock); + InitCriticalSection(StateLock); + InitCriticalSection(DestroyLock); + Randomize; + { Stay below the ranges the systems hand out for outgoing connections - + 32768-60999 on this linux, 49152-65535 on Windows. + And do not rely on Randomize alone to separate runs: on linux it seeds + from the time in whole seconds, so processes started within the same + second - most scenarios take less than one - drew the same ports and + found them still held by the previous process's connections. That was + the cause of the refused connects and bind failures. Mix in the process + id, which differs between consecutive processes. } + PortBase:=20000+((Random(1200)+LongInt(GetProcessID mod 1200)*7) mod 1200)*10; + BuildScenarioTable; + + { ---- child mode: one scenario, then exit ---- } + if (ParamCount>=2) and (ParamStr(1)='--scenario') then + begin + if not SelfTestAccept then + Halt(98); + RunAsChild(StrToIntDef(ParamStr(2),0)); + end; + + { ---- the interrupt-window probe, run directly ---- } + if (ParamCount>0) and (ParamStr(1)='--interrupt-race') then + begin + if not SelfTestAccept then + Halt(98); + TestInterruptOnCompletingRead; + Halt(ScenariosFailed); + end; + + { ---- the crashing opt-in scenario, run directly ---- } + if (ParamCount>0) and (ParamStr(1)='--pump-first') then + begin + if not SelfTestAccept then + Halt(98); + TestPumpFreedBeforeClient; + Halt(ScenariosFailed); + end; + + { ---- runner ---- } + Say('fpwebsocketclient shutdown regression test'); + Say(Format('FPC %s %s-%s', + [{$I %FPCVERSION%},{$I %FPCTARGETCPU%},{$I %FPCTARGETOS%}])); + Say(Format('each scenario runs in its own process, limit %d s', + [ChildLimitMs div 1000])); + Say(''); + + if not SelfTestAccept then + begin + Say('RFC 6455 accept self-test failed - aborting.'); + Halt(98); + end; + + Passed:=0; Failed:=0; Skipped:=0; Hung:=0; Crashed:=0; + For I:=1 to ScenarioCount do + begin + Say(Format('===== scenario %d/%d: %s',[I,ScenarioCount,ScenarioNames[I]])); + V:=RunChild(I,Ex); + Case V of + cvPassed : Inc(Passed); + cvFailed : Inc(Failed); + cvSkipped : Inc(Skipped); + cvHung : + begin + Inc(Hung); + Say(' => HUNG: "'+ScenarioNames[I]+'" did not finish within ' + +IntToStr(ChildLimitMs div 1000)+' s and was killed'); + end; + cvCrashed : + begin + Inc(Crashed); + Say(' => CRASHED: "'+ScenarioNames[I]+'" exit code '+IntToStr(Ex)); + end; + end; + Say(''); + end; + + Say('====================================================='); + Say(Format('passed %d failed %d skipped %d hung %d crashed %d', + [Passed,Failed,Skipped,Hung,Crashed])); + Bad:=Failed+Hung+Crashed; + if Bad=0 then + Say('all executed scenarios passed') + else + Say(Format('%d scenario(s) did not pass',[Bad])); + DoneCriticalSection(StateLock); + DoneCriticalSection(OutLock); + Halt(Bad); +end. base-commit: abcdfa618de1d4d7646698be07fc33b1540a6659 -- 2.51.1