unit Myc.Trade.DataPoint.Impl; interface uses System.SysUtils, Myc.Signals, Myc.Mutable, Myc.Core.Notifier, Myc.Trade.Types, Myc.Trade.DataArray, Myc.Trade.DataPoint; type // 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; public function ProcessData(const Value: S): TState; end; // A generic converter that uses a function reference for the conversion logic. TMycGenericConverter = class(TMycConverter) private FFunc: TConstFunc; protected function ProcessData(const Value: S): TState; override; public 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; protected function ProcessData(const Value: T): TState; override; public 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; // Endpoint that collects data into a series. TMycDataEndpoint = class(TInterfacedObject, TMutable>.IMutable) private FProcessor: TMycContainedProcessor; FTag: TDataProvider.TTag; FDataProvider: TDataProvider; FLookback: Int64; FData: TSeries; FChanged: TEvent; function GetChanged: TSignal; function GetValue: TSeries; function ProcessData(const Value: T): TState; public constructor Create(const ADataProvider: TDataProvider; ALookback: Int64); destructor Destroy; override; end; implementation uses System.TypInfo, 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; begin inherited Create; FSender := TMycDataProvider.Create(Self); end; destructor TMycConverter.Destroy; begin FSender.Free; inherited Destroy; end; function TMycConverter.Broadcast(const Value: T): TState; begin Result := FSender.Broadcast(Value); end; function TMycConverter.GetSender: TDataProvider.IDataProvider; begin Result := FSender; end; { TNullConverter } function TNullConverter.GetSender: TDataProvider.IDataProvider; begin Result := TDataProvider.Null; end; function TNullConverter.ProcessData(const Value: S): TState; begin Result := TState.Null; end; { TMycGenericConverter } constructor TMycGenericConverter.Create(const AFunc: TConstFunc); begin inherited Create; FFunc := AFunc; end; function TMycGenericConverter.ProcessData(const Value: S): TState; begin Result := Broadcast(FFunc(Value)); end; { TMycIndicator } function TMycIndicator.ProcessData(const Value: S): TState; begin Result := Broadcast(Calculate(Value)); end; { TMycDataCounter } constructor TMycDataCounter.Create; begin inherited Create; FCount := 0; end; function TMycDataCounter.ProcessData(const Value: T): TState; begin Result := Broadcast(FCount); inc(FCount); end; { TMycTicker } function TMycTicker.ProcessData(const Values: TArray): TState; begin var done := TLatch.CreateLatch(Length(Values)); // Process each incoming data point for var i := 0 to High(Values) do Broadcast(Values[i]).Signal.Subscribe(done); Result := done.State; end; { TMycRecordFieldReader } constructor TMycRecordFieldReader.Create(const AFieldName: String); begin inherited Create; var Context := TRttiContext.Create; var Field := Context.GetType(TypeInfo(S)).GetField(AFieldName); var TypeT := Context.GetType(TypeInfo(T)); var Fields := Context.GetType(TypeInfo(S)).GetFields; var name := ''; if AFieldName = 'Time' then for var i := 0 to High(Fields) do begin name := name + ' ' + Fields[i].Name; end; Assert(Assigned(Field), 'Field ' + AFieldName + ' not found'); Assert(Field.FieldType.TypeKind = TypeT.TypeKind, 'Incorrect type'); if Assigned(Field) and (Field.FieldType.TypeKind = TypeT.TypeKind) then FOffset := Field.Offset else FOffset := -1; end; function TMycRecordFieldReader.ProcessData(const Values: S): TState; type PT = ^T; begin if FOffset < 0 then exit(TState.Null); var fieldPtr := PByte(@Values); inc(fieldPtr, FOffset); 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; { TMycDataEndpoint } constructor TMycDataEndpoint.Create(const ADataProvider: TDataProvider; ALookback: Int64); begin inherited Create; FDataProvider := ADataProvider; FLookback := ALookback; FProcessor := TMycContainedProcessor.Create(Self, ProcessData); FTag := FDataProvider.Link(FProcessor); end; destructor TMycDataEndpoint.Destroy; begin FDataProvider.Unlink(FTag); FProcessor.Free; inherited; end; function TMycDataEndpoint.GetChanged: TSignal; begin Result := FChanged.Signal; end; function TMycDataEndpoint.GetValue: TSeries; begin Result := FData; end; function TMycDataEndpoint.ProcessData(const Value: T): TState; begin Result := TState.Null; // Add new data point, respecting the lookback period. FData := FData.Add(Value, FLookback); FChanged.Notify; end; end.