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.