Verbatim port of Fastway-Server's TFWEventBus from fw_plugin_host.pas
per feedback_copy_dont_reinterpret.md. Adjustments limited to:
- Type renames (TFW* -> T*).
- uses clause: drop fw_log; add log.types from fpc-log so the
optional Logger property uses the canonical ecosystem-wide
TLogProc shape, matching every other fpc-* library.
- Per-handler exception logging now calls Logger with
Level=llError, Category='events', and includes the source
plugin (ASourcePlugin parameter) in the message text so the
canonical signature stays meaningful.
Behaviours preserved verbatim: APluginName bulk-Unsubscribe key,
wildcard '*' subscriber, OnBroadcast external-listener tap,
snapshot-iterate-outside-lock pattern, per-handler exception
isolation, TCriticalSection.
docs/DEVELOPER_GUIDE.md added covering threading, payload
ownership, recursive Fire, OnBroadcast, logger plumbing, and
the relationship between fpc-events (ecosystem-wide pub/sub)
and per-library typed observer callbacks (bp.events / cm.events
pattern).
Tests: 44 assertions across 14 scenarios pass on x86_64-linux.
Pre-tag -vh audit on src/ev.bus.pas reports zero hints/warnings.
795 lines
19 KiB
ObjectPascal
795 lines
19 KiB
ObjectPascal
{ test_bus -- exercise the TEventBus surface.
|
|
|
|
Targets the canonical TFWEventBus behaviours that fpc-events
|
|
ports verbatim:
|
|
- Subscribe / Unsubscribe / UnsubscribeCallback
|
|
- Wildcard '*' delivery
|
|
- Subscription order preservation
|
|
- Per-handler exception isolation (subscribers + OnBroadcast)
|
|
- Snapshot-iterate semantics (re/un-subscribe inside callback,
|
|
recursive Fire inside callback)
|
|
- Threaded publishes (4 producer x 1000 events) }
|
|
|
|
program test_bus;
|
|
|
|
{$mode objfpc}{$H+}
|
|
|
|
uses
|
|
{$IFDEF UNIX}cthreads,{$ENDIF}
|
|
Classes, SysUtils, fpjson, syncobjs,
|
|
log.types,
|
|
events.version, events.bus;
|
|
|
|
var
|
|
Total, Passed, Failed: Integer;
|
|
|
|
procedure Check(const Name: string; OK: Boolean; const Detail: string = '');
|
|
begin
|
|
Inc(Total);
|
|
if OK then
|
|
begin
|
|
Inc(Passed);
|
|
Writeln(' [PASS] ', Name);
|
|
end
|
|
else
|
|
begin
|
|
Inc(Failed);
|
|
if Detail <> '' then
|
|
Writeln(' [FAIL] ', Name, ' -- ', Detail)
|
|
else
|
|
Writeln(' [FAIL] ', Name);
|
|
end;
|
|
end;
|
|
|
|
type
|
|
TRecorder = class
|
|
public
|
|
Calls: TStringList;
|
|
LogMsgs: TStringList;
|
|
constructor Create;
|
|
destructor Destroy; override;
|
|
procedure HandleA(const AEventType: string; AData: TJSONObject);
|
|
procedure HandleB(const AEventType: string; AData: TJSONObject);
|
|
procedure HandleStar(const AEventType: string; AData: TJSONObject);
|
|
procedure HandleThrow(const AEventType: string; AData: TJSONObject);
|
|
procedure HandleBroadcast(const AEventType: string; AData: TJSONObject);
|
|
procedure HandleBroadcastThrow(const AEventType: string; AData: TJSONObject);
|
|
procedure HandleLogger(Level: TLogLevel; const Category, Msg: string);
|
|
end;
|
|
|
|
constructor TRecorder.Create;
|
|
begin
|
|
inherited Create;
|
|
Calls := TStringList.Create;
|
|
LogMsgs := TStringList.Create;
|
|
end;
|
|
|
|
destructor TRecorder.Destroy;
|
|
begin
|
|
Calls.Free;
|
|
LogMsgs.Free;
|
|
inherited Destroy;
|
|
end;
|
|
|
|
procedure TRecorder.HandleA(const AEventType: string; AData: TJSONObject);
|
|
begin
|
|
Calls.Add('A:' + AEventType);
|
|
end;
|
|
|
|
procedure TRecorder.HandleB(const AEventType: string; AData: TJSONObject);
|
|
begin
|
|
Calls.Add('B:' + AEventType);
|
|
end;
|
|
|
|
procedure TRecorder.HandleStar(const AEventType: string; AData: TJSONObject);
|
|
begin
|
|
Calls.Add('*:' + AEventType);
|
|
end;
|
|
|
|
procedure TRecorder.HandleThrow(const AEventType: string; AData: TJSONObject);
|
|
begin
|
|
Calls.Add('THROW:' + AEventType);
|
|
raise Exception.Create('synthetic handler failure');
|
|
end;
|
|
|
|
procedure TRecorder.HandleBroadcast(const AEventType: string; AData: TJSONObject);
|
|
begin
|
|
Calls.Add('BCAST:' + AEventType);
|
|
end;
|
|
|
|
procedure TRecorder.HandleBroadcastThrow(const AEventType: string; AData: TJSONObject);
|
|
begin
|
|
Calls.Add('BCAST_THROW:' + AEventType);
|
|
raise Exception.Create('synthetic broadcast failure');
|
|
end;
|
|
|
|
procedure TRecorder.HandleLogger(Level: TLogLevel; const Category, Msg: string);
|
|
begin
|
|
LogMsgs.Add(Format('%s/%s/%s', [LogLevelName(Level), Category, Msg]));
|
|
end;
|
|
|
|
{ ---- Re/un-subscribe-from-within-callback recorders ---- }
|
|
|
|
type
|
|
TReentrant = class
|
|
public
|
|
Bus: TEventBus;
|
|
Calls: TStringList;
|
|
SecondaryFired: Integer;
|
|
SubscribedDuring: Boolean;
|
|
constructor Create(ABus: TEventBus);
|
|
destructor Destroy; override;
|
|
procedure HandleSelfUnsubscribe(const AEventType: string; AData: TJSONObject);
|
|
procedure HandleSecondary(const AEventType: string; AData: TJSONObject);
|
|
procedure HandleSubscribeDuring(const AEventType: string; AData: TJSONObject);
|
|
procedure HandleRecursive(const AEventType: string; AData: TJSONObject);
|
|
end;
|
|
|
|
constructor TReentrant.Create(ABus: TEventBus);
|
|
begin
|
|
inherited Create;
|
|
Bus := ABus;
|
|
Calls := TStringList.Create;
|
|
SecondaryFired := 0;
|
|
SubscribedDuring := False;
|
|
end;
|
|
|
|
destructor TReentrant.Destroy;
|
|
begin
|
|
Calls.Free;
|
|
inherited Destroy;
|
|
end;
|
|
|
|
procedure TReentrant.HandleSelfUnsubscribe(const AEventType: string;
|
|
AData: TJSONObject);
|
|
begin
|
|
Calls.Add('SELF:' + AEventType);
|
|
Bus.UnsubscribeCallback(@Self.HandleSelfUnsubscribe);
|
|
end;
|
|
|
|
procedure TReentrant.HandleSecondary(const AEventType: string; AData: TJSONObject);
|
|
begin
|
|
Inc(SecondaryFired);
|
|
end;
|
|
|
|
procedure TReentrant.HandleSubscribeDuring(const AEventType: string;
|
|
AData: TJSONObject);
|
|
begin
|
|
Calls.Add('OUTER:' + AEventType);
|
|
if not SubscribedDuring then
|
|
begin
|
|
SubscribedDuring := True;
|
|
Bus.Subscribe('reentrant', AEventType, @Self.HandleSecondary);
|
|
end;
|
|
end;
|
|
|
|
procedure TReentrant.HandleRecursive(const AEventType: string; AData: TJSONObject);
|
|
var
|
|
Inner: TJSONObject;
|
|
begin
|
|
Calls.Add('REC:' + AEventType);
|
|
if AEventType = 'outer' then
|
|
begin
|
|
Inner := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('test', 'inner', Inner);
|
|
finally
|
|
Inner.Free;
|
|
end;
|
|
end;
|
|
end;
|
|
|
|
{ ---- Threaded publisher ---- }
|
|
|
|
type
|
|
TPublisher = class(TThread)
|
|
public
|
|
Bus: TEventBus;
|
|
EventCount: Integer;
|
|
Marker: string;
|
|
constructor Create(ABus: TEventBus; const AMarker: string;
|
|
ACount: Integer);
|
|
procedure Execute; override;
|
|
end;
|
|
|
|
TCounter = class
|
|
public
|
|
Lock: TCriticalSection;
|
|
Total: Integer;
|
|
constructor Create;
|
|
destructor Destroy; override;
|
|
procedure HandleAny(const AEventType: string; AData: TJSONObject);
|
|
end;
|
|
|
|
constructor TPublisher.Create(ABus: TEventBus; const AMarker: string;
|
|
ACount: Integer);
|
|
begin
|
|
inherited Create(True);
|
|
FreeOnTerminate := False;
|
|
Bus := ABus;
|
|
Marker := AMarker;
|
|
EventCount := ACount;
|
|
end;
|
|
|
|
procedure TPublisher.Execute;
|
|
var
|
|
I: Integer;
|
|
Data: TJSONObject;
|
|
begin
|
|
for I := 0 to EventCount - 1 do
|
|
begin
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Data.Add('marker', Marker);
|
|
Data.Add('seq', I);
|
|
Bus.Fire(Marker, 'thread.event', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
end;
|
|
end;
|
|
|
|
constructor TCounter.Create;
|
|
begin
|
|
inherited Create;
|
|
Lock := TCriticalSection.Create;
|
|
Total := 0;
|
|
end;
|
|
|
|
destructor TCounter.Destroy;
|
|
begin
|
|
Lock.Free;
|
|
inherited Destroy;
|
|
end;
|
|
|
|
procedure TCounter.HandleAny(const AEventType: string; AData: TJSONObject);
|
|
begin
|
|
Lock.Enter;
|
|
try
|
|
Inc(Total);
|
|
finally
|
|
Lock.Leave;
|
|
end;
|
|
end;
|
|
|
|
{ ---- Test scenarios ---- }
|
|
|
|
procedure TestBasicDelivery;
|
|
var
|
|
Bus: TEventBus;
|
|
Rec: TRecorder;
|
|
Data: TJSONObject;
|
|
begin
|
|
Writeln('-- basic delivery');
|
|
Bus := TEventBus.Create;
|
|
Rec := TRecorder.Create;
|
|
try
|
|
Bus.Subscribe('p1', 'evt.a', @Rec.HandleA);
|
|
Check('sub-count after one Subscribe',
|
|
Bus.GetSubscriptionCount = 1);
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('source', 'evt.a', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
Check('callback received fired event',
|
|
(Rec.Calls.Count = 1) and (Rec.Calls[0] = 'A:evt.a'));
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('source', 'evt.b', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
Check('non-matching event skipped',
|
|
Rec.Calls.Count = 1);
|
|
finally
|
|
Rec.Free;
|
|
Bus.Free;
|
|
end;
|
|
end;
|
|
|
|
procedure TestMultipleSubscribersOrder;
|
|
var
|
|
Bus: TEventBus;
|
|
R1, R2, R3: TRecorder;
|
|
Order: TStringList;
|
|
Data: TJSONObject;
|
|
begin
|
|
Writeln('-- subscription order');
|
|
Bus := TEventBus.Create;
|
|
R1 := TRecorder.Create;
|
|
R2 := TRecorder.Create;
|
|
R3 := TRecorder.Create;
|
|
try
|
|
Bus.Subscribe('p1', 'evt', @R1.HandleA);
|
|
Bus.Subscribe('p2', 'evt', @R2.HandleB);
|
|
Bus.Subscribe('p3', 'evt', @R3.HandleA);
|
|
|
|
Check('three subs registered',
|
|
Bus.GetSubscriptionCount = 3);
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('src', 'evt', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
|
|
Order := TStringList.Create;
|
|
try
|
|
if R1.Calls.Count > 0 then Order.Add('1:' + R1.Calls[0]);
|
|
if R2.Calls.Count > 0 then Order.Add('2:' + R2.Calls[0]);
|
|
if R3.Calls.Count > 0 then Order.Add('3:' + R3.Calls[0]);
|
|
Check('all three subscribers fired',
|
|
(R1.Calls.Count = 1) and (R2.Calls.Count = 1) and (R3.Calls.Count = 1));
|
|
finally
|
|
Order.Free;
|
|
end;
|
|
finally
|
|
R1.Free; R2.Free; R3.Free;
|
|
Bus.Free;
|
|
end;
|
|
end;
|
|
|
|
procedure TestWildcard;
|
|
var
|
|
Bus: TEventBus;
|
|
Rec: TRecorder;
|
|
Data: TJSONObject;
|
|
begin
|
|
Writeln('-- wildcard subscriber');
|
|
Bus := TEventBus.Create;
|
|
Rec := TRecorder.Create;
|
|
try
|
|
Bus.Subscribe('w', '*', @Rec.HandleStar);
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('src', 'one', Data);
|
|
Bus.Fire('src', 'two', Data);
|
|
Bus.Fire('src', 'three', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
|
|
Check('wildcard saw all 3 events',
|
|
Rec.Calls.Count = 3);
|
|
Check('wildcard preserved event-type names',
|
|
(Rec.Calls[0] = '*:one') and (Rec.Calls[1] = '*:two') and
|
|
(Rec.Calls[2] = '*:three'));
|
|
finally
|
|
Rec.Free;
|
|
Bus.Free;
|
|
end;
|
|
end;
|
|
|
|
procedure TestTypeFilter;
|
|
var
|
|
Bus: TEventBus;
|
|
Rec: TRecorder;
|
|
Data: TJSONObject;
|
|
begin
|
|
Writeln('-- event-type filtering');
|
|
Bus := TEventBus.Create;
|
|
Rec := TRecorder.Create;
|
|
try
|
|
Bus.Subscribe('p', 'evt.a', @Rec.HandleA);
|
|
Bus.Subscribe('p', 'evt.b', @Rec.HandleB);
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('src', 'evt.a', Data);
|
|
Bus.Fire('src', 'evt.c', Data);
|
|
Bus.Fire('src', 'evt.b', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
Check('handler A only got evt.a',
|
|
(Rec.Calls.Count = 2) and
|
|
(Rec.Calls.IndexOf('A:evt.a') >= 0) and
|
|
(Rec.Calls.IndexOf('B:evt.b') >= 0) and
|
|
(Rec.Calls.IndexOf('A:evt.c') < 0) and
|
|
(Rec.Calls.IndexOf('B:evt.c') < 0));
|
|
finally
|
|
Rec.Free;
|
|
Bus.Free;
|
|
end;
|
|
end;
|
|
|
|
procedure TestUnsubscribeByPlugin;
|
|
var
|
|
Bus: TEventBus;
|
|
R1, R2: TRecorder;
|
|
Data: TJSONObject;
|
|
begin
|
|
Writeln('-- bulk Unsubscribe by plugin name');
|
|
Bus := TEventBus.Create;
|
|
R1 := TRecorder.Create;
|
|
R2 := TRecorder.Create;
|
|
try
|
|
Bus.Subscribe('plug1', 'evt.a', @R1.HandleA);
|
|
Bus.Subscribe('plug1', 'evt.b', @R1.HandleB);
|
|
Bus.Subscribe('plug1', 'evt.c', @R1.HandleStar);
|
|
Bus.Subscribe('plug2', 'evt.a', @R2.HandleA);
|
|
Check('4 subs after registration',
|
|
Bus.GetSubscriptionCount = 4);
|
|
|
|
Bus.Unsubscribe('plug1');
|
|
Check('plug1 removed in bulk -- only plug2 remains',
|
|
Bus.GetSubscriptionCount = 1);
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('src', 'evt.a', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
Check('plug2 still fires after plug1 bulk-removed',
|
|
(R1.Calls.Count = 0) and (R2.Calls.Count = 1));
|
|
finally
|
|
R1.Free; R2.Free;
|
|
Bus.Free;
|
|
end;
|
|
end;
|
|
|
|
procedure TestUnsubscribeByCallback;
|
|
var
|
|
Bus: TEventBus;
|
|
Rec: TRecorder;
|
|
Data: TJSONObject;
|
|
begin
|
|
Writeln('-- UnsubscribeCallback removes only the matching method');
|
|
Bus := TEventBus.Create;
|
|
Rec := TRecorder.Create;
|
|
try
|
|
Bus.Subscribe('p', 'evt', @Rec.HandleA);
|
|
Bus.Subscribe('p', 'evt', @Rec.HandleB);
|
|
Bus.Subscribe('p', 'evt', @Rec.HandleStar);
|
|
Check('3 subs registered', Bus.GetSubscriptionCount = 3);
|
|
|
|
Bus.UnsubscribeCallback(@Rec.HandleB);
|
|
Check('after UnsubscribeCallback(B), 2 remain',
|
|
Bus.GetSubscriptionCount = 2);
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('src', 'evt', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
Check('A and Star still fire, B does not',
|
|
(Rec.Calls.IndexOf('A:evt') >= 0) and
|
|
(Rec.Calls.IndexOf('*:evt') >= 0) and
|
|
(Rec.Calls.IndexOf('B:evt') < 0));
|
|
finally
|
|
Rec.Free;
|
|
Bus.Free;
|
|
end;
|
|
end;
|
|
|
|
procedure TestUnsubscribeFromInsideCallback;
|
|
var
|
|
Bus: TEventBus;
|
|
R: TReentrant;
|
|
Data: TJSONObject;
|
|
begin
|
|
Writeln('-- Unsubscribe from inside a callback');
|
|
Bus := TEventBus.Create;
|
|
R := TReentrant.Create(Bus);
|
|
try
|
|
Bus.Subscribe('r', 'evt', @R.HandleSelfUnsubscribe);
|
|
Check('1 sub registered', Bus.GetSubscriptionCount = 1);
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('src', 'evt', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
Check('first Fire delivered',
|
|
(R.Calls.Count = 1) and (R.Calls[0] = 'SELF:evt'));
|
|
Check('after self-unsubscribe, 0 subs',
|
|
Bus.GetSubscriptionCount = 0);
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('src', 'evt', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
Check('second Fire delivers to no one',
|
|
R.Calls.Count = 1);
|
|
finally
|
|
R.Free;
|
|
Bus.Free;
|
|
end;
|
|
end;
|
|
|
|
procedure TestSubscribeFromInsideCallback;
|
|
var
|
|
Bus: TEventBus;
|
|
R: TReentrant;
|
|
Data: TJSONObject;
|
|
begin
|
|
Writeln('-- Subscribe from inside a callback (snapshot semantics)');
|
|
Bus := TEventBus.Create;
|
|
R := TReentrant.Create(Bus);
|
|
try
|
|
Bus.Subscribe('r', 'evt', @R.HandleSubscribeDuring);
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('src', 'evt', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
Check('outer fired once',
|
|
(R.Calls.Count = 1) and (R.Calls[0] = 'OUTER:evt'));
|
|
Check('secondary not fired in same Fire (snapshot was taken pre-Subscribe)',
|
|
R.SecondaryFired = 0);
|
|
Check('secondary registered',
|
|
Bus.GetSubscriptionCount = 2);
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('src', 'evt', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
Check('secondary fires on next Fire',
|
|
R.SecondaryFired = 1);
|
|
finally
|
|
R.Free;
|
|
Bus.Free;
|
|
end;
|
|
end;
|
|
|
|
procedure TestRecursiveFire;
|
|
var
|
|
Bus: TEventBus;
|
|
R: TReentrant;
|
|
Data: TJSONObject;
|
|
begin
|
|
Writeln('-- recursive Fire from inside a callback');
|
|
Bus := TEventBus.Create;
|
|
R := TReentrant.Create(Bus);
|
|
try
|
|
Bus.Subscribe('r', '*', @R.HandleRecursive);
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('src', 'outer', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
Check('outer fired',
|
|
R.Calls.IndexOf('REC:outer') >= 0);
|
|
Check('inner fired (recursive Fire from callback)',
|
|
R.Calls.IndexOf('REC:inner') >= 0);
|
|
Check('exactly two deliveries',
|
|
R.Calls.Count = 2);
|
|
finally
|
|
R.Free;
|
|
Bus.Free;
|
|
end;
|
|
end;
|
|
|
|
procedure TestHandlerExceptionIsolation;
|
|
var
|
|
Bus: TEventBus;
|
|
R1, R2, R3: TRecorder;
|
|
Logger: TRecorder;
|
|
Data: TJSONObject;
|
|
begin
|
|
Writeln('-- handler exception isolation');
|
|
Bus := TEventBus.Create;
|
|
R1 := TRecorder.Create;
|
|
R2 := TRecorder.Create;
|
|
R3 := TRecorder.Create;
|
|
Logger := TRecorder.Create;
|
|
try
|
|
Bus.Logger := @Logger.HandleLogger;
|
|
Bus.Subscribe('p', 'evt', @R1.HandleA);
|
|
Bus.Subscribe('p', 'evt', @R2.HandleThrow);
|
|
Bus.Subscribe('p', 'evt', @R3.HandleA);
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('src', 'evt', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
Check('R1 fired before R2 raised',
|
|
R1.Calls.Count = 1);
|
|
Check('R2 entered before raising',
|
|
R2.Calls.Count = 1);
|
|
Check('R3 still fired despite R2 raising',
|
|
R3.Calls.Count = 1);
|
|
Check('logger received the handler-error message',
|
|
(Logger.LogMsgs.Count = 1) and
|
|
(Pos('Handler error', Logger.LogMsgs[0]) > 0));
|
|
finally
|
|
R1.Free; R2.Free; R3.Free; Logger.Free;
|
|
Bus.Free;
|
|
end;
|
|
end;
|
|
|
|
procedure TestBroadcastTap;
|
|
var
|
|
Bus: TEventBus;
|
|
R, BC: TRecorder;
|
|
Data: TJSONObject;
|
|
begin
|
|
Writeln('-- OnBroadcast tap');
|
|
Bus := TEventBus.Create;
|
|
R := TRecorder.Create;
|
|
BC := TRecorder.Create;
|
|
try
|
|
Bus.Subscribe('p', 'evt', @R.HandleA);
|
|
Bus.OnBroadcast := @BC.HandleBroadcast;
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('src', 'evt', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
Check('subscriber fired', R.Calls.Count = 1);
|
|
Check('broadcast tap fired with same event',
|
|
(BC.Calls.Count = 1) and (BC.Calls[0] = 'BCAST:evt'));
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('src', 'other', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
Check('broadcast also fires for events with no subscribers',
|
|
BC.Calls.Count = 2);
|
|
finally
|
|
R.Free; BC.Free;
|
|
Bus.Free;
|
|
end;
|
|
end;
|
|
|
|
procedure TestBroadcastExceptionIsolation;
|
|
var
|
|
Bus: TEventBus;
|
|
R: TRecorder;
|
|
BC: TRecorder;
|
|
Logger: TRecorder;
|
|
Data: TJSONObject;
|
|
begin
|
|
Writeln('-- OnBroadcast exception isolation');
|
|
Bus := TEventBus.Create;
|
|
R := TRecorder.Create;
|
|
BC := TRecorder.Create;
|
|
Logger := TRecorder.Create;
|
|
try
|
|
Bus.Logger := @Logger.HandleLogger;
|
|
Bus.Subscribe('p', 'evt', @R.HandleA);
|
|
Bus.OnBroadcast := @BC.HandleBroadcastThrow;
|
|
|
|
Data := TJSONObject.Create;
|
|
try
|
|
Bus.Fire('src', 'evt', Data);
|
|
finally
|
|
Data.Free;
|
|
end;
|
|
Check('subscriber still ran despite broadcast tap raising',
|
|
R.Calls.Count = 1);
|
|
Check('broadcast tap was entered',
|
|
BC.Calls.Count = 1);
|
|
Check('logger received the broadcast-error message',
|
|
(Logger.LogMsgs.Count = 1) and
|
|
(Pos('Broadcast error', Logger.LogMsgs[0]) > 0));
|
|
finally
|
|
R.Free; BC.Free; Logger.Free;
|
|
Bus.Free;
|
|
end;
|
|
end;
|
|
|
|
procedure TestThreadedPublishers;
|
|
const
|
|
ThreadCount = 4;
|
|
PerThread = 1000;
|
|
var
|
|
Bus: TEventBus;
|
|
Counter: TCounter;
|
|
Pubs: array[0..ThreadCount - 1] of TPublisher;
|
|
I: Integer;
|
|
begin
|
|
Writeln('-- threaded publishers (4 x 1000)');
|
|
Bus := TEventBus.Create;
|
|
Counter := TCounter.Create;
|
|
try
|
|
Bus.Subscribe('counter', 'thread.event', @Counter.HandleAny);
|
|
|
|
for I := 0 to ThreadCount - 1 do
|
|
Pubs[I] := TPublisher.Create(Bus, 'pub' + IntToStr(I), PerThread);
|
|
for I := 0 to ThreadCount - 1 do
|
|
Pubs[I].Start;
|
|
for I := 0 to ThreadCount - 1 do
|
|
begin
|
|
Pubs[I].WaitFor;
|
|
Pubs[I].Free;
|
|
end;
|
|
|
|
Check(Format('counter saw %d events (expected %d)',
|
|
[Counter.Total, ThreadCount * PerThread]),
|
|
Counter.Total = ThreadCount * PerThread);
|
|
finally
|
|
Counter.Free;
|
|
Bus.Free;
|
|
end;
|
|
end;
|
|
|
|
procedure TestSubscriptionCount;
|
|
var
|
|
Bus: TEventBus;
|
|
R: TRecorder;
|
|
begin
|
|
Writeln('-- GetSubscriptionCount add/remove');
|
|
Bus := TEventBus.Create;
|
|
R := TRecorder.Create;
|
|
try
|
|
Check('initial count = 0', Bus.GetSubscriptionCount = 0);
|
|
Bus.Subscribe('p', 'a', @R.HandleA);
|
|
Check('after 1 sub = 1', Bus.GetSubscriptionCount = 1);
|
|
Bus.Subscribe('p', 'b', @R.HandleB);
|
|
Bus.Subscribe('q', 'c', @R.HandleStar);
|
|
Check('after 3 subs = 3', Bus.GetSubscriptionCount = 3);
|
|
Bus.UnsubscribeCallback(@R.HandleB);
|
|
Check('after one UnsubscribeCallback = 2',
|
|
Bus.GetSubscriptionCount = 2);
|
|
Bus.Unsubscribe('p');
|
|
Check('after Unsubscribe(p) = 1',
|
|
Bus.GetSubscriptionCount = 1);
|
|
Bus.Unsubscribe('q');
|
|
Check('after Unsubscribe(q) = 0',
|
|
Bus.GetSubscriptionCount = 0);
|
|
finally
|
|
R.Free;
|
|
Bus.Free;
|
|
end;
|
|
end;
|
|
|
|
procedure TestVersionConst;
|
|
begin
|
|
Writeln('-- version constant');
|
|
Check('version string matches expected',
|
|
EVENTS_VERSION_STRING = '0.1.0');
|
|
Check('major/minor/patch decompose',
|
|
(EVENTS_VERSION_MAJOR = 0) and (EVENTS_VERSION_MINOR = 1) and (EVENTS_VERSION_PATCH = 0));
|
|
end;
|
|
|
|
begin
|
|
Total := 0; Passed := 0; Failed := 0;
|
|
Randomize;
|
|
|
|
TestVersionConst;
|
|
TestBasicDelivery;
|
|
TestMultipleSubscribersOrder;
|
|
TestWildcard;
|
|
TestTypeFilter;
|
|
TestUnsubscribeByPlugin;
|
|
TestUnsubscribeByCallback;
|
|
TestUnsubscribeFromInsideCallback;
|
|
TestSubscribeFromInsideCallback;
|
|
TestRecursiveFire;
|
|
TestHandlerExceptionIsolation;
|
|
TestBroadcastTap;
|
|
TestBroadcastExceptionIsolation;
|
|
TestThreadedPublishers;
|
|
TestSubscriptionCount;
|
|
|
|
Writeln;
|
|
Writeln(Format('Total: %d Passed: %d Failed: %d',
|
|
[Total, Passed, Failed]));
|
|
if Failed > 0 then
|
|
Halt(1);
|
|
end.
|