Files
fpc-events/tests/test_bus.pas
T
kenjreno 128ddcb4df v0.1.0: thread-safe pub/sub event bus
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.
2026-05-05 18:13:10 -07:00

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.