diff --git a/Src/Myc.Core.Notifier.pas b/Src/Myc.Core.Notifier.pas index 789ca37..06ff35c 100644 --- a/Src/Myc.Core.Notifier.pas +++ b/Src/Myc.Core.Notifier.pas @@ -19,8 +19,15 @@ type PItem = ^TItem; // Internal structure for storing a receiver and list linkage. TItem = record - Next, Prev: PItem; // Pointers to the next and previous items in the doubly linked list. + private + FNext: PItem; + FPrev: PItem; + function GetNext: PItem; inline; + public + Receiver: T; // The registered interface instance (the event sink). + property Next: PItem read GetNext; + property Prev: PItem read FPrev; end; // Notify function. If this results false, it will be removed from the notify list @@ -31,6 +38,9 @@ type // Bit 0 of the address stores the lock state (0 = locked, 1 = unlocked). [volatile] FList: PItem; + class var + FReverseOnNotify: Boolean; + class constructor CreateClass; // Allocates memory for a new TItem. class function AllocItem: PItem; static; inline; @@ -62,10 +72,18 @@ type // by setting its interface reference to nil. The list item itself is not freed here. procedure Notify(const Func: TNotifyProc); property First: PItem read GetFirst; + + // If this is true, the list is reversed after each Notify, hoping for a better distrubution of events. + class property ReverseOnNotify: Boolean read FReverseOnNotify write FReverseOnNotify; end; implementation +class constructor TMycNotifyList.CreateClass; +begin + FReverseOnNotify := true; +end; + class operator TMycNotifyList.Initialize(out Dest: TMycNotifyList); begin NativeUInt(Dest.FList) := 1; @@ -88,10 +106,10 @@ begin Item := AllocItem; Item.Receiver := Receiver; - Item.Prev := nil; - Item.Next := FList; - if Item.Next <> nil then - Item.Next.Prev := Item; + Item.FPrev := nil; + Item.FNext := FList; + if Item.FNext <> nil then + Item.FNext.FPrev := Item; FList := Item; @@ -133,108 +151,123 @@ var begin Assert(IsLocked); - var rev: PItem := nil; - var tmp: PItem := nil; + // Invariant: Append all detached items at end of the list. We can't simply free them, because the subscription is owned by the client. - while (FList <> nil) and Assigned(FList.Receiver) do + if FReverseOnNotify then begin - Item := FList; - FList := Item.Next; + // Method: stack assigned items and rebuild the list from the stack (reversing the order) - if Func(Item.Receiver) then - begin - Item.Next := rev; - rev := Item; - end - else - begin - // release the receiver - Item.Receiver := nil; + var stack: PItem := nil; + var released: PItem := nil; - Item.Next := tmp; - tmp := Item; + while (FList <> nil) and Assigned(FList.Receiver) do + begin + Item := FList; + FList := Item.FNext; + + if Func(Item.Receiver) then + begin + // item stays valid + Item.FNext := stack; + stack := Item; + end + else + begin + // release the receiver + Item.Receiver := nil; + + Item.FNext := released; + released := Item; + end; end; - end; - while tmp <> nil do - begin - Item := tmp; - tmp := Item.Next; - - Item.Prev := nil; - Item.Next := FList; - FList := Item; - if Item.Next <> nil then - Item.Next.Prev := Item; - end; - - while rev <> nil do - begin - Item := rev; - rev := Item.Next; - - Item.Prev := nil; - Item.Next := FList; - FList := Item; - if Item.Next <> nil then - Item.Next.Prev := Item; - end; - { - tmp := nil; - Last := nil; - Item := FList; - while (Item <> nil) and Assigned(Item.Receiver) do - begin - if not Func(Item.Receiver) then + // The list now only contains old released items. Push the newly released ones. + while released <> nil do begin - // Receiver wants no more notifications, detach it - if Item = FList then - FList := Item.Next; + Item := released; + released := Item.FNext; - if Item.Prev <> nil then - Item.Prev.Next := Item.Next; - if Item.Next <> nil then - Item.Next.Prev := Item.Prev; - - // release the receiver - Item.Receiver := nil; - - // and save the list item in a tmp list for later use - var nxt := Item.Next; - Item.Next := tmp; - tmp := Item; - Item := nxt; - end - else - begin - Last := Item; - Item := Last.Next; - end; - end; - - // Append all detached items at end of the list. We can't simply free them, because the subscription is owned by the client. - while tmp <> nil do - begin - Item := tmp; - tmp := Item.Next; - - if Last = nil then - begin - Item.Prev := nil; - Item.Next := FList; + Item.FPrev := nil; + Item.FNext := FList; FList := Item; - end - else - begin - Item.Prev := Last; - Item.Next := Last.Next; - if Item.Prev <> nil then - Item.Prev.Next := Item; + if Item.FNext <> nil then + Item.FNext.FPrev := Item; + end; + + // Now push the valid items, so that they are at the beginning of the list. + while stack <> nil do + begin + Item := stack; + stack := Item.FNext; + + Item.FPrev := nil; + Item.FNext := FList; + FList := Item; + if Item.FNext <> nil then + Item.FNext.FPrev := Item; + end; + end + else + begin + // Method: filter released items and add then to the end of the list + + var released: PItem := nil; + var lastValid: PItem := nil; + + Item := FList; + while (Item <> nil) and Assigned(Item.Receiver) do + begin + if not Func(Item.Receiver) then + begin + // item is now invalid + if Item = FList then + FList := Item.FNext; + + if Item.FPrev <> nil then + Item.FPrev.FNext := Item.FNext; + if Item.FNext <> nil then + Item.FNext.FPrev := Item.FPrev; + + // release the receiver + Item.Receiver := nil; + + // and save the list item in a tmp list for later use + var nxt := Item.FNext; + Item.FNext := released; + released := Item; + Item := nxt; + end + else + begin + // save a pointer to the last valid item + lastValid := Item; + Item := lastValid.FNext; + end; + end; + + // insert all new released items after the last valid item + while released <> nil do + begin + Item := released; + released := Item.FNext; + + if lastValid = nil then + begin + Item.FPrev := nil; + Item.FNext := FList; + FList := Item; + end + else + begin + Item.FPrev := lastValid; + Item.FNext := lastValid.Next; + if Item.FPrev <> nil then + Item.FPrev.FNext := Item; + end; + if Item.FNext <> nil then + Item.FNext.FPrev := Item; end; - if Item.Next <> nil then - Item.Next.Prev := Item; end; - } end; procedure TMycNotifyList.Release; @@ -270,12 +303,12 @@ begin Item.Receiver := nil; if Item = FList then - FList := Item.Next; + FList := Item.FNext; - if Item.Prev <> nil then - Item.Prev.Next := Item.Next; - if Item.Next <> nil then - Item.Next.Prev := Item.Prev; + if Item.FPrev <> nil then + Item.FPrev.FNext := Item.FNext; + if Item.FNext <> nil then + Item.FNext.FPrev := Item.FPrev; FreeItem(Item); end; @@ -287,4 +320,11 @@ begin Unadvise(TTag(FList)); end; +function TMycNotifyList.TItem.GetNext: PItem; +begin + if not Assigned(Receiver) then + exit(nil); + Result := FNext; +end; + end. diff --git a/Src/Myc.Fmx.Chart.Series.pas b/Src/Myc.Fmx.Chart.Series.pas index 2ac309b..305fb9e 100644 --- a/Src/Myc.Fmx.Chart.Series.pas +++ b/Src/Myc.Fmx.Chart.Series.pas @@ -11,6 +11,7 @@ uses Myc.Trade.Types, Myc.Trade.DataArray, Myc.Trade.DataPoint, + Myc.Trade.DataPoint.Impl, Myc.Fmx.Chart; type diff --git a/Src/Myc.Trade.DataPoint.Impl.pas b/Src/Myc.Trade.DataPoint.Impl.pas index dba38ec..baed245 100644 --- a/Src/Myc.Trade.DataPoint.Impl.pas +++ b/Src/Myc.Trade.DataPoint.Impl.pas @@ -3,12 +3,57 @@ unit Myc.Trade.DataPoint.Impl; interface uses + System.SysUtils, Myc.Signals, Myc.Trade.Types, - Myc.Trade.DataPoint; + Myc.Trade.DataPoint, + Myc.Core.Notifier; type - // Null object implementation for IMycConverter + // Abstract base class for data consumers. + TMycProcessor = class abstract(TInterfacedObject, IMycProcessor) + protected + function ProcessData(const Value: T): TState; virtual; abstract; + end; + + // Concrete data provider that manages a list of processors (listeners). + TMycDataProvider = class abstract(TContainedObject, TDataProvider.IDataProvider) + private + FListeners: TMycNotifyList>; + public + constructor Create(const Controller: IInterface); + destructor Destroy; override; + // Notifies all linked processors. + function Broadcast(const Value: T): TState; + // Link a Processor + function Link(const Processor: IMycProcessor): TDataProvider.TTag; + // Unlink a linked strategy + procedure Unlink(Tag: TDataProvider.TTag); + end; + + // Null object implementation for IDataProvider. + TNullDataProvider = class(TInterfacedObject, TDataProvider.IDataProvider) + public + function Link(const Receiver: IMycProcessor): TDataProvider.TTag; + procedure Unlink(Tag: TDataProvider.TTag); + end; + + // Abstract base class for components that process data of type S and provide data of type T. + TMycConverter = class abstract(TMycProcessor, TConverter.IConverter) + private + FSender: TMycDataProvider; + function GetSender: TDataProvider.IDataProvider; + protected + function ProcessData(const Value: S): TState; override; abstract; + // Broadcasts the given data to all linked processors. + function Broadcast(const Value: T): TState; + public + constructor Create; + destructor Destroy; override; + property Sender: TDataProvider.IDataProvider read GetSender; + end; + + // Null object implementation for IConverter. TNullConverter = class(TInterfacedObject, TConverter.IConverter) private function GetSender: TDataProvider.IDataProvider; @@ -16,22 +61,7 @@ type function ProcessData(const Value: S): TState; end; - TMycConverter = class abstract(TMycProcessor, TConverter.IConverter) - private - FSender: TMycDataProvider; - function GetSender: TDataProvider.IDataProvider; - protected - function ProcessData(const Value: S): TState; override; abstract; - - // Broadcasts the given data to all linked processors. - function Broadcast(const Value: T): TState; - - public - constructor Create; - destructor Destroy; override; - property Sender: TDataProvider.IDataProvider read GetSender; - end; - + // A generic converter that uses a function reference for the conversion logic. TMycGenericConverter = class(TMycConverter) private FFunc: TConstFunc; @@ -41,12 +71,14 @@ type constructor Create(const AFunc: TConstFunc); end; + // A converter specialized for calculating indicators. TMycIndicator = class(TMycConverter) protected function ProcessData(const Value: S): TState; override; final; function Calculate(const Value: S): T; virtual; abstract; end; + // A converter that counts incoming data points and outputs the current count. TMycDataCounter = class(TMycConverter) private FCount: Int64; @@ -56,27 +88,124 @@ type constructor Create; end; + // A converter that takes an array and broadcasts each element individually. TMycTicker = class(TMycConverter, T>) public function ProcessData(const Values: TArray): TState; override; end; + // A converter that reads a specific field from a record using RTTI. TMycRecordFieldReader = class(TMycConverter) private FOffset: Integer; - public constructor Create(const AFieldName: String); function ProcessData(const Values: S): TState; override; end; + // A generic processor implementation that uses a function reference. + TMycGenericProcessor = class(TMycProcessor) + type + TProc = reference to function(const Value: T): TState; + private + FProc: TProc; + protected + function ProcessData(const Value: T): TState; override; final; + public + constructor Create(const AProc: TProc); + end; + + // A processor implementation that is owned by a controller. + TMycContainedProcessor = class(TContainedObject, IMycProcessor) + type + TProc = function(const Value: T): TState of object; + private + FProc: TProc; + function ProcessData(const Value: T): TState; + public + constructor Create(const Controller: IInterface; const AProc: TProc); + end; + implementation uses System.TypInfo, - System.SysUtils, System.RTTI; +{ TMycDataProvider } + +constructor TMycDataProvider.Create(const Controller: IInterface); +begin + inherited Create(Controller); +end; + +destructor TMycDataProvider.Destroy; +begin + FListeners.Finalize; + inherited Destroy; +end; + +function TMycDataProvider.Broadcast(const Value: T): TState; +begin + FListeners.Lock; + try + var item := FListeners.First; + if item = nil then + exit(TState.Null); + + var i := 1; + while item.Next <> nil do + begin + inc(i); + item := item.Next; + end; + + var done := TLatch.CreateLatch(i); + while item <> nil do + begin + item.Receiver.ProcessData(Value).Signal.Subscribe(done); + item := item.Prev; + end; + + Result := done.State; + finally + FListeners.Release; + end; +end; + +function TMycDataProvider.Link(const Processor: IMycProcessor): TDataProvider.TTag; +begin + // Add the Processor to the notification list + FListeners.Lock; + try + Result := FListeners.Advise(Processor); + finally + FListeners.Release; + end; +end; + +procedure TMycDataProvider.Unlink(Tag: TDataProvider.TTag); +begin + FListeners.Lock; + try + FListeners.Unadvise(Tag); + finally + FListeners.Release; + end; +end; + +{ TNullDataProvider } + +function TNullDataProvider.Link(const Receiver: IMycProcessor): TDataProvider.TTag; +begin + Result := nil; +end; + +procedure TNullDataProvider.Unlink(Tag: TDataProvider.TTag); +begin + // Do nothing in the null implementation. +end; + { TMycConverter } constructor TMycConverter.Create; @@ -101,7 +230,7 @@ begin Result := FSender; end; -{ TNullConverter } +{ TNullConverter } function TNullConverter.GetSender: TDataProvider.IDataProvider; begin @@ -133,6 +262,8 @@ begin Result := Broadcast(Calculate(Value)); end; +{ TMycDataCounter } + constructor TMycDataCounter.Create; begin inherited Create; @@ -145,6 +276,8 @@ begin inc(FCount); end; +{ TMycTicker } + function TMycTicker.ProcessData(const Values: TArray): TState; begin var done := TLatch.CreateLatch(Length(Values)); @@ -156,6 +289,8 @@ begin Result := done.State; end; +{ TMycRecordFieldReader } + constructor TMycRecordFieldReader.Create(const AFieldName: String); begin inherited Create; @@ -193,4 +328,30 @@ begin Result := Broadcast(PT(fieldPtr)^); end; +{ TMycGenericProcessor } + +constructor TMycGenericProcessor.Create(const AProc: TProc); +begin + inherited Create; + FProc := AProc; +end; + +function TMycGenericProcessor.ProcessData(const Value: T): TState; +begin + Result := FProc(Value); +end; + +{ TMycContainedProcessor } + +constructor TMycContainedProcessor.Create(const Controller: IInterface; const AProc: TProc); +begin + inherited Create(Controller); + FProc := AProc; +end; + +function TMycContainedProcessor.ProcessData(const Value: T): TState; +begin + Result := FProc(Value); +end; + end. diff --git a/Src/Myc.Trade.DataPoint.pas b/Src/Myc.Trade.DataPoint.pas index fed9483..dd4bce7 100644 --- a/Src/Myc.Trade.DataPoint.pas +++ b/Src/Myc.Trade.DataPoint.pas @@ -4,8 +4,7 @@ interface uses Myc.Signals, - Myc.Trade.Types, - Myc.Core.Notifier; + Myc.Trade.Types; type // Represents a time-stamped data point in a series. @@ -19,11 +18,6 @@ type function ProcessData(const Value: T): TState; end; - TMycProcessor = class abstract(TInterfacedObject, IMycProcessor) - protected - function ProcessData(const Value: T): TState; virtual; abstract; - end; - TDataProvider = record type TTag = Pointer; @@ -56,48 +50,6 @@ type class property Null: IDataProvider read FNull; end; - TMycDataProvider = class abstract(TContainedObject, TDataProvider.IDataProvider) - private - FListeners: TMycNotifyList>; - public - constructor Create(const Controller: IInterface); - destructor Destroy; override; - // Notifies all linked processors. - function Broadcast(const Value: T): TState; - procedure Notify(const Func: TMycNotifyList>.TNotifyProc); - // Link a Processor - function Link(const Processor: IMycProcessor): TDataProvider.TTag; - // Unlink a linked strategy - procedure Unlink(Tag: TDataProvider.TTag); - end; - - TMycGenericProcessor = class(TMycProcessor) - type - TProc = reference to function(const Value: T): TState; - private - FProc: TProc; - protected - function ProcessData(const Value: T): TState; override; final; - public - constructor Create(const AProc: TProc); - end; - - TMycContainedProcessor = class(TContainedObject, IMycProcessor) - type - TProc = function(const Value: T): TState of object; - private - FProc: TProc; - function ProcessData(const Value: T): TState; - public - constructor Create(const Controller: IInterface; const AProc: TProc); - end; - - TNullDataProvider = class(TInterfacedObject, TDataProvider.IDataProvider) - public - function Link(const Receiver: IMycProcessor): TDataProvider.TTag; - procedure Unlink(Tag: TDataProvider.TTag); - end; - // Interface helper for IMycConverter providing the null object pattern. TConverter = record type @@ -150,18 +102,6 @@ implementation uses Myc.Trade.DataPoint.Impl; -{ TNullDataProvider } - -function TNullDataProvider.Link(const Receiver: IMycProcessor): TDataProvider.TTag; -begin - Result := nil; -end; - -procedure TNullDataProvider.Unlink(Tag: TDataProvider.TTag); -begin - // Do nothing in the null implementation. -end; - { TDataPoint } constructor TDataPoint.Create(ATime: TDateTime; const AData: T); @@ -216,104 +156,6 @@ begin FDataProvider.Unlink(Tag); end; -{ TMycGenericProcessor } - -constructor TMycGenericProcessor.Create(const AProc: TProc); -begin - inherited Create; - FProc := AProc; -end; - -function TMycGenericProcessor.ProcessData(const Value: T): TState; -begin - Result := FProc(Value); -end; - -{ TMycContainedProcessor } - -constructor TMycContainedProcessor.Create(const Controller: IInterface; const AProc: TProc); -begin - inherited Create(Controller); - FProc := AProc; -end; - -function TMycContainedProcessor.ProcessData(const Value: T): TState; -begin - Result := FProc(Value); -end; - -{ TMycDataProvider } - -constructor TMycDataProvider.Create(const Controller: IInterface); -begin - inherited Create(Controller); -end; - -destructor TMycDataProvider.Destroy; -begin - FListeners.Finalize; - inherited Destroy; -end; - -function TMycDataProvider.Broadcast(const Value: T): TState; -begin - FListeners.Lock; - try - var item := FListeners.First; - if item = nil then - exit(TState.Null); - - var i := 1; - while item.Next <> nil do - begin - inc(i); - item := item.Next; - end; - - var done := TLatch.CreateLatch(i); - while item <> nil do - begin - item.Receiver.ProcessData(Value).Signal.Subscribe(done); - item := item.Prev; - end; - - Result := done.State; - finally - FListeners.Release; - end; -end; - -procedure TMycDataProvider.Notify(const Func: TMycNotifyList>.TNotifyProc); -begin - FListeners.Lock; - try - FListeners.Notify(Func); - finally - FListeners.Release; - end; -end; - -function TMycDataProvider.Link(const Processor: IMycProcessor): TDataProvider.TTag; -begin - // Add the Processor to the notification list - FListeners.Lock; - try - Result := FListeners.Advise(Processor); - finally - FListeners.Release; - end; -end; - -procedure TMycDataProvider.Unlink(Tag: TDataProvider.TTag); -begin - FListeners.Lock; - try - FListeners.Unadvise(Tag); - finally - FListeners.Release; - end; -end; - { TConverter } class constructor TConverter.CreateClass;