235 lines
6.3 KiB
ObjectPascal
235 lines
6.3 KiB
ObjectPascal
unit Myc.Signals.FMX;
|
|
|
|
interface
|
|
|
|
uses
|
|
System.Classes,
|
|
System.SysUtils,
|
|
Myc.Signals;
|
|
|
|
type
|
|
TSignalComponentHelper = class helper for TComponent
|
|
type
|
|
TMsgProc = reference to procedure(out IsDone: Boolean);
|
|
public
|
|
function ProcessSignal(const Signal: TSignal; const Proc: TMsgProc): TComponent; overload;
|
|
function ProcessSignal(const Signal: TSignal; const Proc: TProc): TComponent; overload;
|
|
function ProcessSignal(const Signal: TSignal; const Done: TState; const Proc: TProc): TComponent; overload;
|
|
end;
|
|
|
|
TSignalSyncHelper = record helper for TSignal
|
|
function Queue(const Proc: TProc; Delay: Integer = 0): TSignal.TSubscription; overload;
|
|
function Queue(Thread: TThread; const Proc: TProc; Delay: Integer = 0): TSignal.TSubscription; overload;
|
|
end;
|
|
|
|
implementation
|
|
|
|
uses
|
|
System.Diagnostics,
|
|
System.Messaging,
|
|
FMX.Types;
|
|
|
|
type
|
|
TSignalSubscriber = class(TComponent, TSignal.ISubscriber)
|
|
private
|
|
FSignal: TSignal;
|
|
FSigSubscr: TSignal.TSubscription;
|
|
FProc: TSignalComponentHelper.TMsgProc;
|
|
FNotified: Integer;
|
|
FIdleSubscrId: TMessageSubscriptionId;
|
|
FNext, FPrev: TSignalSubscriber;
|
|
class var
|
|
FQueued: Integer;
|
|
FCount: Integer;
|
|
FFirst: TSignalSubscriber;
|
|
FCurr: TSignalSubscriber;
|
|
class procedure HandleSignals(Timeout: Int64);
|
|
public
|
|
constructor Create(AOwner: TComponent; const ASignal: TSignal; const AProc: TSignalComponentHelper.TMsgProc); reintroduce;
|
|
destructor Destroy; override;
|
|
procedure AfterConstruction; override;
|
|
procedure BeforeDestruction; override;
|
|
function Notify: Boolean;
|
|
end;
|
|
|
|
TSyncSubscriber = class(TInterfacedObject, TSignal.ISubscriber)
|
|
Thread: TThread;
|
|
Delay: Integer;
|
|
Timestamp: Int64;
|
|
Proc: TProc;
|
|
class var
|
|
Timer: TStopwatch;
|
|
function Notify: Boolean;
|
|
class constructor CreateClass;
|
|
end;
|
|
|
|
class constructor TSyncSubscriber.CreateClass;
|
|
begin
|
|
Timer := TStopwatch.StartNew;
|
|
end;
|
|
|
|
function TSyncSubscriber.Notify: Boolean;
|
|
begin
|
|
if Assigned(Proc) then
|
|
begin
|
|
var cProc := Proc;
|
|
Proc := nil;
|
|
TThread.ForceQueue(Thread, procedure begin cProc() end, Delay - (Timer.ElapsedMilliseconds - Timestamp));
|
|
end;
|
|
Result := false;
|
|
end;
|
|
|
|
{ TSignalSubscription }
|
|
|
|
constructor TSignalSubscriber.Create(AOwner: TComponent; const ASignal: TSignal; const AProc: TSignalComponentHelper.TMsgProc);
|
|
begin
|
|
inherited Create(AOwner);
|
|
FSignal := ASignal;
|
|
FProc := AProc;
|
|
|
|
if FFirst = nil then
|
|
begin
|
|
FIdleSubscrId :=
|
|
TMessageManager
|
|
.DefaultManager
|
|
.SubscribeToMessage(TIdleMessage, procedure(const Sender: TObject; const M: TMessage) begin HandleSignals(50); end);
|
|
end;
|
|
|
|
FPrev := nil;
|
|
if FFirst <> nil then
|
|
FNext := FFirst;
|
|
if FNext <> nil then
|
|
FNext.FPrev := Self;
|
|
FFirst := Self;
|
|
|
|
FSigSubscr := FSignal.Subscribe(Self);
|
|
end;
|
|
|
|
destructor TSignalSubscriber.Destroy;
|
|
begin
|
|
FSigSubscr.Unsubscribe;
|
|
|
|
if FFirst = Self then
|
|
FFirst := FNext;
|
|
if FNext <> nil then
|
|
FNext.FPrev := FPrev;
|
|
if FPrev <> nil then
|
|
FPrev.FNext := FNext;
|
|
|
|
if FFirst = nil then
|
|
begin
|
|
TMessageManager.DefaultManager.Unsubscribe(TIdleMessage, FIdleSubscrId);
|
|
end;
|
|
|
|
inherited;
|
|
end;
|
|
|
|
procedure TSignalSubscriber.AfterConstruction;
|
|
begin
|
|
inherited;
|
|
Notify;
|
|
end;
|
|
|
|
procedure TSignalSubscriber.BeforeDestruction;
|
|
begin
|
|
FSigSubscr.Unsubscribe;
|
|
if FCurr = Self then
|
|
FCurr := FCurr.FNext;
|
|
inherited;
|
|
end;
|
|
|
|
class procedure TSignalSubscriber.HandleSignals(Timeout: Int64);
|
|
begin
|
|
if AtomicExchange(FQueued, 0) = 0 then
|
|
exit;
|
|
|
|
var Stopwatch := TStopwatch.StartNew;
|
|
var doBreak := false;
|
|
|
|
if FCurr = nil then
|
|
FCurr := FFirst;
|
|
|
|
while FCurr <> nil do
|
|
begin
|
|
var sub := FCurr;
|
|
FCurr := sub.FNext;
|
|
|
|
if AtomicExchange(sub.FNotified, 0) = 1 then
|
|
begin
|
|
if not (csDestroying in sub.ComponentState) then
|
|
begin
|
|
var done := false;
|
|
try
|
|
sub.FProc(done);
|
|
finally
|
|
if done then
|
|
sub.Free;
|
|
end;
|
|
end;
|
|
|
|
doBreak := AtomicDecrement(FCount) = 1;
|
|
end;
|
|
|
|
doBreak := doBreak or (Stopwatch.ElapsedMilliseconds > Timeout);
|
|
|
|
if doBreak then
|
|
begin
|
|
AtomicExchange(FQueued, 1);
|
|
break;
|
|
end;
|
|
end;
|
|
end;
|
|
|
|
function TSignalSubscriber.Notify: Boolean;
|
|
begin
|
|
if AtomicExchange(FNotified, 1) = 0 then
|
|
AtomicIncrement(FCount);
|
|
AtomicExchange(FQueued, 1);
|
|
Result := true;
|
|
end;
|
|
|
|
function TSignalComponentHelper.ProcessSignal(const Signal: TSignal; const Proc: TProc): TComponent;
|
|
begin
|
|
var cProc: TProc := Proc;
|
|
Result := ProcessSignal(Signal, procedure(out IsDone: Boolean) begin cProc(); end);
|
|
end;
|
|
|
|
function TSignalComponentHelper.ProcessSignal(const Signal: TSignal; const Proc: TMsgProc): TComponent;
|
|
begin
|
|
Result := TSignalSubscriber.Create(Self, Signal, Proc);
|
|
end;
|
|
|
|
function TSignalComponentHelper.ProcessSignal(const Signal: TSignal; const Done: TState; const Proc: TProc): TComponent;
|
|
begin
|
|
var cProc: TProc := Proc;
|
|
var cDone: TState := Done;
|
|
Result :=
|
|
ProcessSignal(
|
|
Signal,
|
|
procedure(out IsDone: Boolean)
|
|
begin
|
|
cProc();
|
|
IsDone := cDone.IsSet;
|
|
end
|
|
);
|
|
end;
|
|
|
|
{ TSignalSyncHelper }
|
|
|
|
function TSignalSyncHelper.Queue(const Proc: TProc; Delay: Integer = 0): TSignal.TSubscription;
|
|
begin
|
|
Result := Queue(nil, Proc, Delay);
|
|
end;
|
|
|
|
function TSignalSyncHelper.Queue(Thread: TThread; const Proc: TProc; Delay: Integer = 0): TSignal.TSubscription;
|
|
begin
|
|
var Subscr := TSyncSubscriber.Create;
|
|
Subscr.Thread := Thread;
|
|
Subscr.Proc := Proc;
|
|
Subscr.Delay := Delay;
|
|
Subscr.Timestamp := TSyncSubscriber.Timer.ElapsedMilliseconds;
|
|
Result := Subscribe(Subscr);
|
|
end;
|
|
|
|
end.
|