diff --git a/AuraTrader/AuraTrader.dproj b/AuraTrader/AuraTrader.dproj index 18a82b0..0ded5a4 100644 --- a/AuraTrader/AuraTrader.dproj +++ b/AuraTrader/AuraTrader.dproj @@ -4,7 +4,7 @@ 20.3 FMX True - Release + Debug Win64 AuraTrader 3 diff --git a/AuraTrader/MainForm.pas b/AuraTrader/MainForm.pas index 6d8de84..303241f 100644 --- a/AuraTrader/MainForm.pas +++ b/AuraTrader/MainForm.pas @@ -105,13 +105,13 @@ type FApplication: IAuraApplication; FModulesItem: TTreeViewItem; function SelectedSymbol: String; - procedure ExecuteStrategy(const Symbol: String; Timeframe: TTimeframe; const Processor: IMycProcessor>); + procedure ExecuteStrategy(const Symbol: String; Timeframe: TTimeframe; const Consumer: IConsumer>); public procedure NewWorkspace; function CurrLayout: T; procedure AlignControl(Control: TControl); - function CreateStrategy2(Timeframe: TTimeframe): IMycProcessor>; + function CreateStrategy2(Timeframe: TTimeframe): IConsumer>; published property OnEvent: TNotifyEvent read FOnEvent write FOnEvent; @@ -285,7 +285,7 @@ begin Control.Align := TAlignLayout.Top; end; -function TForm1.CreateStrategy2(Timeframe: TTimeframe): IMycProcessor>; +function TForm1.CreateStrategy2(Timeframe: TTimeframe): IConsumer>; type TSignal = record Sig: Double; @@ -297,9 +297,9 @@ var panel: TMycChart.TPanel; begin var ticker := TConverter.CreateIdentity>; - Result := ticker; + Result := ticker.Consumer; - var OhlcPoint := ticker.DataProvider.Chain>(TConverter.CreateOhlcAggregation(Timeframe)); + var OhlcPoint := ticker.Producer.Chain>(TConverter.CreateOhlcAggregation(Timeframe)); var Ohlc := OhlcPoint.Field('Data'); @@ -426,7 +426,7 @@ begin conv.Chain(Signal); - var pnl := Signal.DataProvider.Field('pnl'); + var pnl := Signal.Producer.Field('pnl'); var FEquity: Double := 10000; var FInit: Boolean := false; @@ -470,21 +470,21 @@ begin panel.AddOhlcSeries(Ohlc); panel.AddDoubleSeries(Hull, TAlphaColors.Cornflowerblue, 2); panel.AddDoubleSeries(Sma, TAlphaColors.Brown, 1.5); - panel.AddDoubleSeries(Signal.DataProvider.Field('Entry'), TAlphaColors.Green, 1); - panel.AddDoubleSeries(Signal.DataProvider.Field('SL'), TAlphaColors.Red, 2); + panel.AddDoubleSeries(Signal.Producer.Field('Entry'), TAlphaColors.Green, 1); + panel.AddDoubleSeries(Signal.Producer.Field('SL'), TAlphaColors.Red, 2); var mean := TConverter, Double>.CreateGeneric(TIndicators.CreateMean()); TConverter.Join([Hull, Sma]).Chain(mean); - panel.AddDoubleSeries(mean.DataProvider, TAlphaColors.Blue, 5); + panel.AddDoubleSeries(mean.Producer, TAlphaColors.Blue, 5); var pnlChart := TMycChart.Create(Self); AlignControl(pnlChart); pnlChart.Height := Layout.ChildrenRect.Width * 9 / 24; pnlChart.Lookback.Value := 50000; - pnlChart.SetXAxisCounter(equity.DataProvider); + pnlChart.SetXAxisCounter(equity.Producer); //////////// var EMAFactory := TEMA.Create; @@ -498,7 +498,7 @@ begin var equityEMA := equity - .DataProvider + .Producer .Chain(TConverter.FieldToRecord(EMAFactory.Input, 'Price')) .Chain(EMAConv) .Chain(TConverter.FieldOfRecord(EMAFactory.Output, 'MA')); @@ -506,7 +506,7 @@ begin ////////////// panel := pnlChart.AddPanel; - panel.AddDoubleSeries(equity.DataProvider, TAlphaColors.Blue, 3); + panel.AddDoubleSeries(equity.Producer, TAlphaColors.Blue, 3); panel.AddDoubleSeries(equityEMA, TAlphaColors.Gray, 2); ///// @@ -533,7 +533,7 @@ begin Result := Res; end; -procedure TForm1.ExecuteStrategy(const Symbol: String; Timeframe: TTimeframe; const Processor: IMycProcessor>); +procedure TForm1.ExecuteStrategy(const Symbol: String; Timeframe: TTimeframe; const Consumer: IConsumer>); begin var terminated := TFlag.CreateObserver(FTerminate.Signal).State; @@ -550,9 +550,9 @@ begin ); var OhlcPoint := lastPrice.Chain>(TConverter.CreateTickAggregation(Timeframe)); - OhlcPoint.DataProvider.Link(Processor); + OhlcPoint.Producer.Link(Consumer); - var dataProvider := + var Producer := TConverter>, TArray>>.CreateGeneric( function(const Values: TArray>): TArray> begin @@ -566,16 +566,16 @@ begin end ); - dataProvider.DataProvider.Link(ticker); + Producer.Producer.Link(ticker); - FProcessDone := FProcessDone + FServer.ProcessData(Symbol, terminated, dataProvider); + FProcessDone := FProcessDone + FServer.ProcessData(Symbol, terminated, Producer); {$else} var ticker := TConverter.CreateTicker>; - ticker.DataProvider.Chain(Processor); + ticker.Producer.Chain(Consumer); - FProcessDone := FProcessDone + FServer.ProcessData(Symbol, terminated, ticker); + FProcessDone := FProcessDone + FServer.ProcessData(Symbol, terminated, ticker.Consumer); {$endif} end; @@ -609,7 +609,7 @@ begin var ticker := TConverter.CreateIdentity>; - var OhlcPoint := ticker.DataProvider.Chain>(TConverter.CreateOhlcAggregation(timeframe)); + var OhlcPoint := ticker.Producer.Chain>(TConverter.CreateOhlcAggregation(timeframe)); // var OhlcTicker := TConverter.CreateIdentity>; // var OhlcPoint := OhlcTicker.Sender; @@ -674,7 +674,7 @@ begin } ///// - ExecuteStrategy(Symbol, timeframe, ticker); + ExecuteStrategy(Symbol, timeframe, ticker.Consumer); end; procedure TForm1.Strat2ButtonClick(Sender: TObject); @@ -687,7 +687,7 @@ begin //ExecuteStrategy(Symbol, timeframe, CreateStrategy2(timeframe)); var tstStrat := StrategyTest.CreateStrategy1(timeframe); - ExecuteStrategy(Symbol, timeframe, tstStrat); + ExecuteStrategy(Symbol, timeframe, tstStrat.Consumer); var Layout := CurrLayout; if Layout = nil then @@ -698,10 +698,10 @@ begin pnlChart.Height := Layout.ChildrenRect.Width * 9 / 24; pnlChart.Lookback.Value := 50000; - pnlChart.SetXAxisCounter(tstStrat.DataProvider); + pnlChart.SetXAxisCounter(tstStrat.Producer); var panel := pnlChart.AddPanel; - panel.AddDoubleSeries(tstStrat.DataProvider, TAlphaColors.Blue, 3); + panel.AddDoubleSeries(tstStrat.Producer, TAlphaColors.Blue, 3); end; end. diff --git a/AuraTrader/StrategyTest.pas b/AuraTrader/StrategyTest.pas index d82bc80..cff102b 100644 --- a/AuraTrader/StrategyTest.pas +++ b/AuraTrader/StrategyTest.pas @@ -30,7 +30,7 @@ type begin var ticker := TConverter.CreateIdentity>; - var OhlcPoint := ticker.DataProvider.Chain>(TConverter.CreateOhlcAggregation(Timeframe)); + var OhlcPoint := ticker.Producer.Chain>(TConverter.CreateOhlcAggregation(Timeframe)); var Ohlc := OhlcPoint.Field('Data'); @@ -183,7 +183,7 @@ begin ) ); - Result := TConverter, Double>.Construct(ticker, equity); + Result := TConverter, Double>.Construct(ticker.Consumer, equity); end; end. diff --git a/Src/Myc.Fmx.Chart.Series.pas b/Src/Myc.Fmx.Chart.Series.pas index 4182e57..cb527ae 100644 --- a/Src/Myc.Fmx.Chart.Series.pas +++ b/Src/Myc.Fmx.Chart.Series.pas @@ -18,7 +18,7 @@ type { TChartCustomLayer } TChartCustomLayer = class abstract(TMycChart.TDataLayer) private - FDataProvider: TDataProvider; + FProducer: TProducer; FReceiver: TLazy>; FData: TSeries; protected @@ -33,7 +33,7 @@ type const YForm: TFunc ); override; abstract; public - constructor Create(AParent: TMycChart.TPanel; const ADataProvider: TDataProvider); + constructor Create(AParent: TMycChart.TPanel; const AProducer: TProducer); property Data: TSeries read FData; end; @@ -53,7 +53,7 @@ type public constructor Create( AParent: TMycChart.TPanel; - const ADataProvider: TDataProvider; + const AProducer: TProducer; const AUpColor, ADownColor: TMutable ); end; @@ -74,7 +74,7 @@ type public constructor Create( AParent: TMycChart.TPanel; - const ADataProvider: TDataProvider; + const AProducer: TProducer; const ALineColor: TAlphaColor; ALineWidth: Single ); @@ -83,7 +83,7 @@ type { TChartXAxisLayer } TChartXAxisLayer = class(TMycChart.TXAxisLayer) private - FDataProvider: TDataProvider; + FProducer: TProducer; FReceiver: TLazy>; FData: TSeries; protected @@ -92,7 +92,7 @@ type function GetCount: Int64; override; function GetTotalCount: Int64; override; public - constructor Create(AOwner: TMycChart; const ADataProvider: TDataProvider); + constructor Create(AOwner: TMycChart; const AProducer: TProducer); property Data: TSeries read FData; end; @@ -104,7 +104,7 @@ type // Provide a formatted timestamp string for a given data index. function GetCaption(Idx: Int64): String; override; final; public - constructor Create(AOwner: TMycChart; ATimeframe: TTimeframe; const ADataProvider: TDataProvider); + constructor Create(AOwner: TMycChart; ATimeframe: TTimeframe; const AProducer: TProducer); end; implementation @@ -118,12 +118,12 @@ uses constructor TChartLineLayer.Create( AParent: TMycChart.TPanel; - const ADataProvider: TDataProvider; + const AProducer: TProducer; const ALineColor: TAlphaColor; ALineWidth: Single ); begin - inherited Create(AParent, ADataProvider); + inherited Create(AParent, AProducer); FLineColor := ALineColor; FLineWidth := ALineWidth; end; @@ -194,11 +194,11 @@ end; constructor TChartOhlcLayer.Create( AParent: TMycChart.TPanel; - const ADataProvider: TDataProvider; + const AProducer: TProducer; const AUpColor, ADownColor: TMutable ); begin - inherited Create(AParent, ADataProvider); + inherited Create(AParent, AProducer); FUpColor := AUpColor; FDownColor := ADownColor; @@ -267,12 +267,12 @@ end; { TChartCustomLayer } -constructor TChartCustomLayer.Create(AParent: TMycChart.TPanel; const ADataProvider: TDataProvider); +constructor TChartCustomLayer.Create(AParent: TMycChart.TPanel; const AProducer: TProducer); begin inherited Create(AParent); - FDataProvider := ADataProvider; + FProducer := AProducer; - FReceiver := FDataProvider.CreateEndpoint(Owner.Lookback.Value); + FReceiver := FProducer.CreateEndpoint(Owner.Lookback.Value); end; function TChartCustomLayer.GetCount: Int64; @@ -292,11 +292,11 @@ end; { TChartXAxisLayer } -constructor TChartXAxisLayer.Create(AOwner: TMycChart; const ADataProvider: TDataProvider); +constructor TChartXAxisLayer.Create(AOwner: TMycChart; const AProducer: TProducer); begin inherited Create(AOwner); - FDataProvider := ADataProvider; - FReceiver := FDataProvider.CreateEndpoint(Owner.Lookback.Value); + FProducer := AProducer; + FReceiver := FProducer.CreateEndpoint(Owner.Lookback.Value); end; function TChartXAxisLayer.GetCaption(Idx: Int64): String; @@ -321,9 +321,9 @@ end; { TChartXAxisTimestampLayer } -constructor TChartXAxisTimestampLayer.Create(AOwner: TMycChart; ATimeframe: TTimeframe; const ADataProvider: TDataProvider); +constructor TChartXAxisTimestampLayer.Create(AOwner: TMycChart; ATimeframe: TTimeframe; const AProducer: TProducer); begin - inherited Create(AOwner, ADataProvider); + inherited Create(AOwner, AProducer); FTimeframe := ATimeframe; end; diff --git a/Src/Myc.Fmx.Chart.pas b/Src/Myc.Fmx.Chart.pas index ab334cf..8484cbb 100644 --- a/Src/Myc.Fmx.Chart.pas +++ b/Src/Myc.Fmx.Chart.pas @@ -97,13 +97,13 @@ type destructor Destroy; override; function AddOhlcSeries( - const DataProvider: TDataProvider; + const Producer: TProducer; const AUpColor: TAlphaColor = TAlphaColors.Green; const ADownColor: TAlphaColor = TAlphaColors.Red ): TMycChart.TDataLayer; function AddDoubleSeries( - const DataProvider: TDataProvider; + const Producer: TProducer; const ALineColor: TAlphaColor = TAlphaColors.Cornflowerblue; const ALineWidth: Single = 1.5 ): TMycChart.TDataLayer; @@ -150,8 +150,8 @@ type function AddPanel: TPanel; // Sets the master series that defines the time scale (X-axis). - function SetXAxisSeries(Timeframe: TTimeframe; const DataProvider: TDataProvider): TMycChart.TXAxisLayer; overload; - function SetXAxisCounter(const DataProvider: TDataProvider): TMycChart.TXAxisLayer; overload; + function SetXAxisSeries(Timeframe: TTimeframe; const Producer: TProducer): TMycChart.TXAxisLayer; overload; + function SetXAxisCounter(const Producer: TProducer): TMycChart.TXAxisLayer; overload; property Lookback: TWriteable read FLookback write FLookback; property NeedRepaint: TFlag read FNeedRepaint; @@ -632,19 +632,19 @@ begin Repaint; end; -function TMycChart.SetXAxisCounter(const DataProvider: TDataProvider): TMycChart.TXAxisLayer; +function TMycChart.SetXAxisCounter(const Producer: TProducer): TMycChart.TXAxisLayer; begin FXAxisSeries.Free; - FXAxisSeries := TChartXAxisLayer.Create(Self, DataProvider.Chain(TConverter.CreateCounter)); + FXAxisSeries := TChartXAxisLayer.Create(Self, Producer.Chain(TConverter.CreateCounter)); Result := FXAxisSeries; end; -function TMycChart.SetXAxisSeries(Timeframe: TTimeframe; const DataProvider: TDataProvider): TMycChart.TXAxisLayer; +function TMycChart.SetXAxisSeries(Timeframe: TTimeframe; const Producer: TProducer): TMycChart.TXAxisLayer; begin FXAxisSeries.Free; - FXAxisSeries := TChartXAxisTimestampLayer.Create(Self, Timeframe, DataProvider); + FXAxisSeries := TChartXAxisTimestampLayer.Create(Self, Timeframe, Producer); Result := FXAxisSeries; end; @@ -670,23 +670,22 @@ begin end; function TMycChart.TPanel.AddDoubleSeries( - const DataProvider: TDataProvider; + const Producer: TProducer; const ALineColor: TAlphaColor = TAlphaColors.Cornflowerblue; const ALineWidth: Single = 1.5 ): TMycChart.TDataLayer; begin - Result := TChartLineLayer.Create(Self, DataProvider, ALineColor, ALineWidth); + Result := TChartLineLayer.Create(Self, Producer, ALineColor, ALineWidth); FLayers.Add(Result); end; function TMycChart.TPanel.AddOhlcSeries( - const DataProvider: TDataProvider; + const Producer: TProducer; const AUpColor: TAlphaColor = TAlphaColors.Green; const ADownColor: TAlphaColor = TAlphaColors.Red ): TMycChart.TDataLayer; begin - Result := - TChartOhlcLayer.Create(Self, DataProvider, TMutable.Constant(AUpColor), TMutable.Constant(ADownColor)); + Result := TChartOhlcLayer.Create(Self, Producer, TMutable.Constant(AUpColor), TMutable.Constant(ADownColor)); FLayers.Add(Result); end; diff --git a/Src/Myc.Trade.DataPoint.Impl.pas b/Src/Myc.Trade.DataPoint.Impl.pas index 085af95..93a1533 100644 --- a/Src/Myc.Trade.DataPoint.Impl.pas +++ b/Src/Myc.Trade.DataPoint.Impl.pas @@ -15,68 +15,91 @@ uses type // Abstract base class for data consumers. - TMycProcessor = class abstract(TInterfacedObject, IMycProcessor) + TMycConsumer = class abstract(TContainedObject, IConsumer) + strict private + type + TNull = class(TInterfacedObject, IConsumer) + protected + function Consume(const Value: T): TState; + end; + class var + FNull: IConsumer; + class constructor CreateClass; protected - function Execute(const Value: T; const Proc: TConstFunc): TState; - function ProcessData(const Value: T): TState; virtual; abstract; - end; - - // Concrete data provider that manages a list of processors (listeners). - TMycContainedDataProvider = class abstract(TContainedObject, IDataProvider) - private - FListeners: TMycNotifyList>; + function Consume(const Value: T): TState; virtual; abstract; 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): TTag; - // Unlink a linked strategy - procedure Unlink(Tag: TTag); + class property Null: IConsumer read FNull; end; - TMycSequence = class(TMycProcessor, IMycDataSequence) + TMycProducer = class(TInterfacedObject, IProducer) + strict private + type + TNull = class(TInterfacedObject, IProducer) + public + function Link(const Consumer: IConsumer): TTag; + procedure Unlink(Tag: TTag); + end; + class var + FNull: IProducer; + class constructor CreateClass; private - FDataProviders: TArray>; - function GetCount: Integer; - function GetDataProvider(Idx: Integer): TDataProvider; - protected - function ProcessData(const Value: T): TState; override; - function ProcessDataProvider(Idx: Integer; const Value: T): TState; - public - constructor Create(ACount: Integer); - destructor Destroy; override; - end; - - // Null object implementation for IDataProvider. - TNullDataProvider = class(TInterfacedObject, IDataProvider) - public - function Link(const Receiver: IMycProcessor): TTag; - procedure Unlink(Tag: TTag); - end; - - // Abstract base class for components that process data of type S and provide data of type T. - TMycConverter = class abstract(TMycProcessor, IConverter) - private - FDataProvider: TMycContainedDataProvider; - function GetDataProvider: IDataProvider; - protected - function ProcessData(const Value: S): TState; override; abstract; - // Broadcasts the given data to all linked processors. - function Broadcast(const Value: T): TState; + FListeners: TMycNotifyList>; public constructor Create; destructor Destroy; override; - property DataProvider: IDataProvider read GetDataProvider; + function Broadcast(const Value: T): TState; + function Link(const Consumer: IConsumer): TTag; + procedure Unlink(Tag: TTag); + class property Null: IProducer read FNull; end; - // Null object implementation for IConverter. - TNullConverter = class(TInterfacedObject, IConverter) + // Concrete producer that manages a list of consumers (listeners). + TMycContainedProducer = class(TContainedObject, IProducer) private - function GetDataProvider: IDataProvider; + FListeners: TMycNotifyList>; public - function ProcessData(const Value: S): TState; + constructor Create(const Controller: IInterface); + destructor Destroy; override; + function Broadcast(const Value: T): TState; + function Link(const Consumer: IConsumer): TTag; + procedure Unlink(Tag: TTag); + end; + + // Abstract base class for components that now act as a producer and contain a consumer. + TMycConverter = class abstract(TMycProducer, IConverter) + strict private + type + TNull = class(TInterfacedObject, IConverter) + private + function GetConsumer: IConsumer; + function Link(const Consumer: IConsumer): TTag; + procedure Unlink(Tag: TTag); + end; + class var + FNull: IConverter; + class constructor CreateClass; + private + type + TConsumer = class(TMycConsumer) + private + FOwner: TMycConverter; + protected + function Consume(const Value: S): TState; override; final; + public + constructor Create(AOwner: TMycConverter); + end; + + private + FConsumer: TConsumer; + function GetConsumer: IConsumer; + protected + // To be implemented by descendants to perform the actual conversion. + function Consume(const Value: S): TState; virtual; abstract; + public + constructor Create; + destructor Destroy; override; + + class property Null: IConverter read FNull; end; // A generic converter that uses a function reference for the conversion logic. @@ -84,7 +107,7 @@ type private FFunc: TConstFunc; protected - function ProcessData(const Value: S): TState; override; + function Consume(const Value: S): TState; override; public constructor Create(const AFunc: TConstFunc); end; @@ -93,7 +116,7 @@ type private FFunc: TConstFunc.TBroadcastProc, TState>; protected - function ProcessData(const Value: S): TState; override; + function Consume(const Value: S): TState; override; public constructor Create(const AFunc: TConstFunc.TBroadcastProc, TState>); end; @@ -103,21 +126,14 @@ type FFunc: TConstFunc; FQueue: TState; protected - function ProcessData(const Value: S): TState; override; + function Consume(const Value: S): TState; override; public constructor Create(const AFunc: TConstFunc); end; TMycIdentityConverter = class(TMycConverter) protected - function ProcessData(const Value: T): TState; override; final; - 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; + function Consume(const Value: T): TState; override; final; end; // A converter that counts incoming data points and outputs the current count. @@ -125,47 +141,40 @@ type private FCount: Int64; protected - function ProcessData(const Value: T): TState; override; + function Consume(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; + protected + function Consume(const Values: TArray): TState; override; end; // A converter that reads a specific field from a record using RTTI. TMycRecordFieldReader = class(TMycConverter) private FOffset: Integer; + protected + function Consume(const Values: S): TState; override; 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) + // A consumer implementation that is owned by a controller. + TMycGenericConsumer = class(TMycConsumer) private FProc: TConstFunc; protected - function ProcessData(const Value: T): TState; override; final; - public - constructor Create(const AProc: TConstFunc); - end; - - // A processor implementation that is owned by a controller. - TMycContainedProcessor = class(TContainedObject, IMycProcessor) - private - FProc: TConstFunc; - function ProcessData(const Value: T): TState; + function Consume(const Value: T): TState; override; public constructor Create(const Controller: IInterface; const AProc: TConstFunc); end; // Endpoint that collects data into a series. TMycDataEndpoint = class(TInterfacedObject, TLazy>.ILazy) + private type PItem = ^TItem; TItem = record @@ -173,18 +182,18 @@ type Value: T; end; private - FProcessor: TMycContainedProcessor; + FConsumer: TMycGenericConsumer; FTag: TTag; - FDataProvider: IDataProvider; + FProducer: IProducer; FLookback: Int64; FChanged: TFlag; FLock: TLightweightMREW; FFirst: PItem; FCount: Integer; function GetChanged: TState; - function ProcessData(const Value: T): TState; + function Consume(const Value: T): TState; public - constructor Create(const ADataProvider: TDataProvider; ALookback: Int64); + constructor Create(const AProducer: TProducer; ALookback: Int64); destructor Destroy; override; function Update(var Value: TSeries): Boolean; end; @@ -196,9 +205,10 @@ type function GetBarStartTime(const TimeStamp: TDateTime; const Timeframe: TTimeframe): TDateTime; function GetCurrentBar: TDataPoint; function GetTimeframe: TTimeframe; + protected + function Consume(const Value: TDataPoint): TState; override; public constructor Create(const ATimeframe: TTimeframe); - function ProcessData(const Value: TDataPoint): TState; override; property CurrentBar: TDataPoint read GetCurrentBar; property Timeframe: TTimeframe read GetTimeframe; end; @@ -207,7 +217,7 @@ type private FQueue: TState; protected - function ProcessData(const Value: T): TState; override; final; + function Consume(const Value: T): TState; override; final; end; TOhlcAggregation = class(TMycConverter, TDataPoint>) @@ -217,42 +227,58 @@ type function GetBarStartTime(const TimeStamp: TDateTime; const Timeframe: TTimeframe): TDateTime; function GetCurrentBar: TDataPoint; function GetTimeframe: TTimeframe; + protected + function Consume(const Value: TDataPoint): TState; override; public constructor Create(const ATimeframe: TTimeframe); - function ProcessData(const Value: TDataPoint): TState; override; property CurrentBar: TDataPoint read GetCurrentBar; property Timeframe: TTimeframe read GetTimeframe; end; // Endpoint that collects data into a series. - TMycDataJoin = class(TInterfacedObject, IDataProvider>) + TMycDataJoin = class(TInterfacedObject, IProducer>) private - FReceivers: array of record - DataProvider: IDataProvider; - Processor: TMycContainedProcessor; + FConsumers: array of record + Producer: IProducer; + Consumer: TMycGenericConsumer; Tag: TTag; Queue: TQueue; end; - FDataProvider: TMycContainedDataProvider>; + FContainedProvider: TMycContainedProducer>; FLock: TSpinLock; FCount: Integer; - function ProcessData(Idx: Integer; const Value: T): TState; + function Consume(Idx: Integer; const Value: T): TState; + function Link(const Consumer: IConsumer>): TTag; + procedure Unlink(Tag: TTag); public - constructor Create(const ADataProviders: TArray>); + constructor Create(const AProducers: TArray>); destructor Destroy; override; - property DataProvider: TMycContainedDataProvider> read FDataProvider implements IDataProvider>; end; - TMycComposedConverter = class(TMycProcessor, IConverter) + TMycComposedConverter = class(TInterfacedObject, IConverter) private - FProcessor: IMycProcessor; - FDataProvider: TDataProvider; + FConsumer: IConsumer; + FProducer: TProducer; protected - function ProcessData(const Value: S): TState; override; - function GetDataProvider: IDataProvider; + function GetConsumer: IConsumer; + function Link(const Consumer: IConsumer): TTag; + procedure Unlink(Tag: TTag); public - constructor Create(const AProcessor: IMycProcessor; const ADataProvider: TDataProvider); + constructor Create(const AConsumer: IConsumer; const AProducer: TProducer); + end; + + TMycSequence = class(TInterfacedObject, IConsumer) + private + FProducers: TArray>; + protected + function Consume(const Value: T): TState; + function ProcessProducer(Idx: Integer; const Value: T): TState; + public + constructor Create(ACount: Integer); + destructor Destroy; override; + + class function CreateSequence(Count: Integer; out Producers: TArray>): IConsumer; static; end; implementation @@ -265,28 +291,25 @@ uses Winapi.Windows, Myc.TaskManager; -{ TMycProcessor } - -function TMycProcessor.Execute(const Value: T; const Proc: TConstFunc): TState; +class constructor TMycConsumer.CreateClass; begin - var cValue := Value; - Result := TaskManager.RunTask(function: TState begin Result := Proc(cValue); end); + FNull := TNull.Create; end; -{ TMycContainedDataProvider } +{ TMycContainedProducer } -constructor TMycContainedDataProvider.Create(const Controller: IInterface); +constructor TMycContainedProducer.Create(const Controller: IInterface); begin inherited Create(Controller); end; -destructor TMycContainedDataProvider.Destroy; +destructor TMycContainedProducer.Destroy; begin FListeners.Finalize; inherited Destroy; end; -function TMycContainedDataProvider.Broadcast(const Value: T): TState; +function TMycContainedProducer.Broadcast(const Value: T): TState; begin FListeners.Lock; try @@ -304,7 +327,7 @@ begin var done := TLatch.CreateLatch(i); while item <> nil do begin - item.Receiver.ProcessData(Value).Signal.Subscribe(done); + item.Receiver.Consume(Value).Signal.Subscribe(done); item := item.Prev; end; @@ -314,18 +337,18 @@ begin end; end; -function TMycContainedDataProvider.Link(const Processor: IMycProcessor): TTag; +function TMycContainedProducer.Link(const Consumer: IConsumer): TTag; begin - // Add the Processor to the notification list + // Add the Consumer to the notification list FListeners.Lock; try - Result := FListeners.Advise(Processor); + Result := FListeners.Advise(Consumer); finally FListeners.Release; end; end; -procedure TMycContainedDataProvider.Unlink(Tag: TTag); +procedure TMycContainedProducer.Unlink(Tag: TTag); begin FListeners.Lock; try @@ -335,52 +358,47 @@ begin end; end; -{ TNullDataProvider } - -function TNullDataProvider.Link(const Receiver: IMycProcessor): TTag; -begin - Result := nil; -end; - -procedure TNullDataProvider.Unlink(Tag: TTag); -begin - // Do nothing in the null implementation. -end; - { TMycConverter } constructor TMycConverter.Create; begin inherited Create; - FDataProvider := TMycContainedDataProvider.Create(Self); + FConsumer := TConsumer.Create(Self); +end; + +{ TMycConverter } + +class constructor TMycConverter.CreateClass; +begin + FNull := TNull.Create; end; destructor TMycConverter.Destroy; begin - FDataProvider.Free; + FConsumer.Free; inherited Destroy; end; -function TMycConverter.Broadcast(const Value: T): TState; +function TMycConverter.GetConsumer: IConsumer; begin - Result := FDataProvider.Broadcast(Value); + Result := FConsumer; end; -function TMycConverter.GetDataProvider: IDataProvider; +{ TMycConverter.TNull } + +function TMycConverter.TNull.GetConsumer: IConsumer; begin - Result := FDataProvider; + Result := TMycConsumer.Null; end; -{ TNullConverter } - -function TNullConverter.GetDataProvider: IDataProvider; +function TMycConverter.TNull.Link(const Consumer: IConsumer): TTag; begin - Result := TDataProvider.Null; + Result := nil; end; -function TNullConverter.ProcessData(const Value: S): TState; +procedure TMycConverter.TNull.Unlink(Tag: TTag); begin - Result := TState.Null; + // NOP end; { TMycGenericConverter } @@ -391,18 +409,11 @@ begin FFunc := AFunc; end; -function TMycGenericConverter.ProcessData(const Value: S): TState; +function TMycGenericConverter.Consume(const Value: S): TState; begin Result := Broadcast(FFunc(Value)); end; -{ TMycIndicator } - -function TMycIndicator.ProcessData(const Value: S): TState; -begin - Result := Execute(Value, function(const Value: S): TState begin Result := Broadcast(Calculate(Value)); end); -end; - { TMycDataCounter } constructor TMycDataCounter.Create; @@ -411,7 +422,7 @@ begin FCount := 0; end; -function TMycDataCounter.ProcessData(const Value: T): TState; +function TMycDataCounter.Consume(const Value: T): TState; begin Result := Broadcast(FCount); inc(FCount); @@ -419,11 +430,10 @@ end; { TMycTicker } -function TMycTicker.ProcessData(const Values: TArray): TState; +function TMycTicker.Consume(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); @@ -440,8 +450,6 @@ begin var Field := Context.GetType(TypeInfo(S)).GetField(AFieldName); var TypeT := Context.GetType(TypeInfo(T)); - var Fields := Context.GetType(TypeInfo(S)).GetFields; - Assert(Assigned(Field), 'Field ' + AFieldName + ' not found'); Assert(Field.FieldType.TypeKind = TypeT.TypeKind, 'Incorrect type'); @@ -451,7 +459,7 @@ begin FOffset := -1; end; -function TMycRecordFieldReader.ProcessData(const Values: S): TState; +function TMycRecordFieldReader.Consume(const Values: S): TState; type PT = ^T; begin @@ -463,48 +471,35 @@ begin Result := Broadcast(PT(fieldPtr)^); end; -{ TMycGenericProcessor } +{ TMycGenericConsumer } -constructor TMycGenericProcessor.Create(const AProc: TConstFunc); -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: TConstFunc); +constructor TMycGenericConsumer.Create(const Controller: IInterface; const AProc: TConstFunc); begin inherited Create(Controller); FProc := AProc; end; -function TMycContainedProcessor.ProcessData(const Value: T): TState; +function TMycGenericConsumer.Consume(const Value: T): TState; begin Result := FProc(Value); end; { TMycDataEndpoint } -constructor TMycDataEndpoint.Create(const ADataProvider: TDataProvider; ALookback: Int64); +constructor TMycDataEndpoint.Create(const AProducer: TProducer; ALookback: Int64); begin inherited Create; - FDataProvider := ADataProvider; + FProducer := AProducer; FLookback := ALookback; - FProcessor := TMycContainedProcessor.Create(Self, ProcessData); - FTag := FDataProvider.Link(FProcessor); + FConsumer := TMycGenericConsumer.Create(Self, Consume); + FTag := FProducer.Link(FConsumer); end; destructor TMycDataEndpoint.Destroy; begin - FDataProvider.Unlink(FTag); - FProcessor.Free; + FProducer.Unlink(FTag); + FConsumer.Free; inherited; end; @@ -513,7 +508,7 @@ begin Result := FChanged.State end; -function TMycDataEndpoint.ProcessData(const Value: T): TState; +function TMycDataEndpoint.Consume(const Value: T): TState; begin FLock.BeginWrite; try @@ -573,15 +568,13 @@ function TTickAggregation.GetBarStartTime(const TimeStamp: TDateTime; const Time var baseTime: TDateTime; begin - // Align the time grid to UTC 0:00 using functions from System.DateUtils + // Implementation is unchanged baseTime := RecodeMilliSecond(TimeStamp, 0); - case Timeframe of S: Result := baseTime; S5: Result := RecodeSecond(baseTime, SecondOf(TimeStamp) - (SecondOf(TimeStamp) mod 5)); S15: Result := RecodeSecond(baseTime, SecondOf(TimeStamp) - (SecondOf(TimeStamp) mod 15)); S30: Result := RecodeSecond(baseTime, SecondOf(TimeStamp) - (SecondOf(TimeStamp) mod 30)); - M: Result := RecodeSecond(baseTime, 0); M2: Result := RecodeMinute(RecodeSecond(baseTime, 0), MinuteOf(TimeStamp) - (MinuteOf(TimeStamp) mod 2)); M3: Result := RecodeMinute(RecodeSecond(baseTime, 0), MinuteOf(TimeStamp) - (MinuteOf(TimeStamp) mod 3)); @@ -589,30 +582,21 @@ begin M10: Result := RecodeMinute(RecodeSecond(baseTime, 0), MinuteOf(TimeStamp) - (MinuteOf(TimeStamp) mod 10)); M15: Result := RecodeMinute(RecodeSecond(baseTime, 0), MinuteOf(TimeStamp) - (MinuteOf(TimeStamp) mod 15)); M30: Result := RecodeMinute(RecodeSecond(baseTime, 0), MinuteOf(TimeStamp) - (MinuteOf(TimeStamp) mod 30)); - H: Result := RecodeMinute(RecodeSecond(baseTime, 0), 0); H2: Result := RecodeHour(RecodeMinute(RecodeSecond(baseTime, 0), 0), HourOf(TimeStamp) - (HourOf(TimeStamp) mod 2)); H3: Result := RecodeHour(RecodeMinute(RecodeSecond(baseTime, 0), 0), HourOf(TimeStamp) - (HourOf(TimeStamp) mod 3)); H4: Result := RecodeHour(RecodeMinute(RecodeSecond(baseTime, 0), 0), HourOf(TimeStamp) - (HourOf(TimeStamp) mod 4)); H8: Result := RecodeHour(RecodeMinute(RecodeSecond(baseTime, 0), 0), HourOf(TimeStamp) - (HourOf(TimeStamp) mod 8)); H12: Result := RecodeHour(RecodeMinute(RecodeSecond(baseTime, 0), 0), HourOf(TimeStamp) - (HourOf(TimeStamp) mod 12)); - D: Result := StartOfTheDay(TimeStamp); - // D2, D3 are uncommon; this is a simple modulo-based approach relative to TDateTime's epoch. D2: Result := Floor(TimeStamp) - (Floor(TimeStamp) mod 2); D3: Result := Floor(TimeStamp) - (Floor(TimeStamp) mod 3); - W: Result := TimeStamp.StartOfTheWeek; - MN: Result := TimeStamp.StartOfTheMonth; - // Quarter alignment MN3: Result := RecodeMonth(TimeStamp.StartOfTheMonth, (MonthOf(TimeStamp) - 1) div 3 * 3 + 1); - // Half-year alignment MN6: Result := RecodeMonth(TimeStamp.StartOfTheMonth, (MonthOf(TimeStamp) - 1) div 6 * 6 + 1); - Y: Result := TimeStamp.StartOfTheYear; else - // Fallback for any undefined timeframe Result := 0; end; end; @@ -627,24 +611,21 @@ begin Result := FTimeframe; end; -function TTickAggregation.ProcessData(const Value: TDataPoint): TState; +function TTickAggregation.Consume(const Value: TDataPoint): TState; var barStartTime: TDateTime; lastBarTime: TDateTime; begin - // Update bar for the strategy's timeframe barStartTime := GetBarStartTime(Value.Time, FTimeframe); lastBarTime := FCurrentBar.Time; if (barStartTime > lastBarTime) then begin - // A new bar starts, so the previous one is now complete. if (lastBarTime > 0) then begin Result := Broadcast(FCurrentBar); end; - // Start a new bar, Volume is 1 because this is the first tick. FCurrentBar.Data.Open := Value.Data; FCurrentBar.Data.High := Value.Data; FCurrentBar.Data.Low := Value.Data; @@ -654,13 +635,11 @@ begin end else begin - // Update the currently aggregating bar if Value.Data > FCurrentBar.Data.High then FCurrentBar.Data.High := Value.Data; if Value.Data < FCurrentBar.Data.Low then FCurrentBar.Data.Low := Value.Data; FCurrentBar.Data.Close := Value.Data; - // Volume is the number of ticks needed to build the complete bar. FCurrentBar.Data.Volume := FCurrentBar.Data.Volume + 1; end; end; @@ -671,44 +650,44 @@ constructor TMycSequence.Create(ACount: Integer); begin inherited Create; - SetLength(FDataProviders, ACount); - for var i := 0 to High(FDataProviders) do - FDataProviders[i] := TMycContainedDataProvider.Create(Self); + SetLength(FProducers, ACount); + for var i := 0 to High(FProducers) do + FProducers[i] := TMycContainedProducer.Create(Self); end; destructor TMycSequence.Destroy; begin - for var i := High(FDataProviders) downto 0 do - FDataProviders[i].Free; + for var i := High(FProducers) downto 0 do + FProducers[i].Free; inherited; end; -function TMycSequence.GetCount: Integer; +class function TMycSequence.CreateSequence(Count: Integer; out Producers: TArray>): IConsumer; +var + Seq: TMycSequence; begin - Result := Length(FDataProviders); + Seq := TMycSequence.Create(Count); + Result := Seq; + + SetLength(Producers, Length(Seq.FProducers)); + for var i := 0 to High(Producers) do + Producers[i] := Seq.FProducers[i]; end; -function TMycSequence.GetDataProvider(Idx: Integer): TDataProvider; +function TMycSequence.Consume(const Value: T): TState; begin - Result := FDataProviders[Idx]; + Result := ProcessProducer(0, Value); end; -function TMycSequence.ProcessData(const Value: T): TState; +function TMycSequence.ProcessProducer(Idx: Integer; const Value: T): TState; begin - Result := ProcessDataProvider(0, Value); -end; - -function TMycSequence.ProcessDataProvider(Idx: Integer; const Value: T): TState; -begin - if Idx >= Length(FDataProviders) then + if Idx >= Length(FProducers) then exit; - Result := - TaskManager - .RunTask(FDataProviders[idx].Broadcast(Value), function: TState begin Result := ProcessDataProvider(1 + idx, Value); end); + Result := TaskManager.RunTask(FProducers[idx].Broadcast(Value), function: TState begin Result := ProcessProducer(1 + idx, Value); end); end; -function TMycIdentityConverter.ProcessData(const Value: T): TState; +function TMycIdentityConverter.Consume(const Value: T): TState; begin Result := Broadcast(Value); end; @@ -721,7 +700,7 @@ begin FFunc := AFunc; end; -function TMycGenericParallelConverter.ProcessData(const Value: S): TState; +function TMycGenericParallelConverter.Consume(const Value: S): TState; begin var cValue := Value; Result := TaskManager.RunTask(FQueue, function: TState begin Result := Broadcast(FFunc(cValue)); end); @@ -729,7 +708,9 @@ begin FQueue := Result; end; -function TMycParallelConverter.ProcessData(const Value: T): TState; +{ TMycParallelConverter } + +function TMycParallelConverter.Consume(const Value: T): TState; begin var cValue := Value; Result := TaskManager.RunTask(FQueue, function: TState begin Result := Broadcast(cValue); end); @@ -745,18 +726,14 @@ begin end; function TOhlcAggregation.GetBarStartTime(const TimeStamp: TDateTime; const Timeframe: TTimeframe): TDateTime; -var - baseTime: TDateTime; begin - // Align the time grid to UTC 0:00 using functions from System.DateUtils - baseTime := RecodeMilliSecond(TimeStamp, 0); - + // Same implementation as TTickAggregation + var baseTime := RecodeMilliSecond(TimeStamp, 0); case Timeframe of S: Result := baseTime; S5: Result := RecodeSecond(baseTime, SecondOf(TimeStamp) - (SecondOf(TimeStamp) mod 5)); S15: Result := RecodeSecond(baseTime, SecondOf(TimeStamp) - (SecondOf(TimeStamp) mod 15)); S30: Result := RecodeSecond(baseTime, SecondOf(TimeStamp) - (SecondOf(TimeStamp) mod 30)); - M: Result := RecodeSecond(baseTime, 0); M2: Result := RecodeMinute(RecodeSecond(baseTime, 0), MinuteOf(TimeStamp) - (MinuteOf(TimeStamp) mod 2)); M3: Result := RecodeMinute(RecodeSecond(baseTime, 0), MinuteOf(TimeStamp) - (MinuteOf(TimeStamp) mod 3)); @@ -764,30 +741,21 @@ begin M10: Result := RecodeMinute(RecodeSecond(baseTime, 0), MinuteOf(TimeStamp) - (MinuteOf(TimeStamp) mod 10)); M15: Result := RecodeMinute(RecodeSecond(baseTime, 0), MinuteOf(TimeStamp) - (MinuteOf(TimeStamp) mod 15)); M30: Result := RecodeMinute(RecodeSecond(baseTime, 0), MinuteOf(TimeStamp) - (MinuteOf(TimeStamp) mod 30)); - H: Result := RecodeMinute(RecodeSecond(baseTime, 0), 0); H2: Result := RecodeHour(RecodeMinute(RecodeSecond(baseTime, 0), 0), HourOf(TimeStamp) - (HourOf(TimeStamp) mod 2)); H3: Result := RecodeHour(RecodeMinute(RecodeSecond(baseTime, 0), 0), HourOf(TimeStamp) - (HourOf(TimeStamp) mod 3)); H4: Result := RecodeHour(RecodeMinute(RecodeSecond(baseTime, 0), 0), HourOf(TimeStamp) - (HourOf(TimeStamp) mod 4)); H8: Result := RecodeHour(RecodeMinute(RecodeSecond(baseTime, 0), 0), HourOf(TimeStamp) - (HourOf(TimeStamp) mod 8)); H12: Result := RecodeHour(RecodeMinute(RecodeSecond(baseTime, 0), 0), HourOf(TimeStamp) - (HourOf(TimeStamp) mod 12)); - D: Result := StartOfTheDay(TimeStamp); - // D2, D3 are uncommon; this is a simple modulo-based approach relative to TDateTime's epoch. D2: Result := Floor(TimeStamp) - (Floor(TimeStamp) mod 2); D3: Result := Floor(TimeStamp) - (Floor(TimeStamp) mod 3); - W: Result := TimeStamp.StartOfTheWeek; - MN: Result := TimeStamp.StartOfTheMonth; - // Quarter alignment MN3: Result := RecodeMonth(TimeStamp.StartOfTheMonth, (MonthOf(TimeStamp) - 1) div 3 * 3 + 1); - // Half-year alignment MN6: Result := RecodeMonth(TimeStamp.StartOfTheMonth, (MonthOf(TimeStamp) - 1) div 6 * 6 + 1); - Y: Result := TimeStamp.StartOfTheYear; else - // Fallback for any undefined timeframe Result := 0; end; end; @@ -802,141 +770,250 @@ begin Result := FTimeframe; end; -function TOhlcAggregation.ProcessData(const Value: TDataPoint): TState; +function TOhlcAggregation.Consume(const Value: TDataPoint): TState; var barStartTime: TDateTime; lastBarTime: TDateTime; begin - // Update bar for the strategy's timeframe barStartTime := GetBarStartTime(Value.Time, FTimeframe); lastBarTime := FCurrentBar.Time; if (barStartTime > lastBarTime) then begin - // A new bar starts, so the previous one is now complete. if (lastBarTime > 0) then begin Result := Broadcast(FCurrentBar); end; - // Start a new bar, Volume is 1 because this is the first tick. FCurrentBar.Data := Value.Data; FCurrentBar.Time := barStartTime; end else begin - // Update the currently aggregating bar if Value.Data.High > FCurrentBar.Data.High then FCurrentBar.Data.High := Value.Data.High; if Value.Data.Low < FCurrentBar.Data.Low then FCurrentBar.Data.Low := Value.Data.Low; FCurrentBar.Data.Close := Value.Data.Close; - // Volume is the number of ticks needed to build the complete bar. FCurrentBar.Data.Volume := FCurrentBar.Data.Volume + Value.Data.Volume; end; end; { TMycDataJoin } -constructor TMycDataJoin.Create(const ADataProviders: TArray>); +constructor TMycDataJoin.Create(const AProducers: TArray>); begin inherited Create; - SetLength(FReceivers, Length(ADataProviders)); + SetLength(FConsumers, Length(AProducers)); FLock := TSpinLock.Create(false); - FCount := Length(FReceivers); + FCount := Length(FConsumers); - FDataProvider := TMycContainedDataProvider>.Create(Self); + FContainedProvider := TMycContainedProducer>.Create(Self); var cFunc := function(Idx: Integer): TConstFunc begin - Result := function(const Value: T): TState begin ProcessData(Idx, Value); end + Result := function(const Value: T): TState begin Result := Consume(Idx, Value); end end; - for var i := 0 to High(FReceivers) do - with FReceivers[i] do + for var i := 0 to High(FConsumers) do + with FConsumers[i] do begin Queue := TQueue.Create; - DataProvider := ADataProviders[i]; - Processor := TMycContainedProcessor.Create(Self, cFunc(i)); - Tag := DataProvider.Link(Processor); + Producer := AProducers[i]; + Consumer := TMycGenericConsumer.Create(Self, cFunc(i)); + Tag := Producer.Link(Consumer); end; end; destructor TMycDataJoin.Destroy; begin - for var i := High(FReceivers) downto 0 do - with FReceivers[i] do + for var i := High(FConsumers) downto 0 do + with FConsumers[i] do begin - DataProvider.Unlink(FReceivers[i].Tag); - Processor.Free; + Producer.Unlink(FConsumers[i].Tag); + Consumer.Free; Queue.Free; end; - FDataProvider.Free; + FContainedProvider.Free; inherited; end; -function TMycDataJoin.ProcessData(Idx: Integer; const Value: T): TState; +function TMycDataJoin.Link(const Consumer: IConsumer>): TTag; +begin + Result := FContainedProvider.Link(Consumer); +end; + +procedure TMycDataJoin.Unlink(Tag: TTag); +begin + FContainedProvider.Unlink(Tag); +end; + +function TMycDataJoin.Consume(Idx: Integer; const Value: T): TState; begin var Arr: TArray; FLock.Enter; try - if FReceivers[Idx].Queue.Count = 0 then + if FConsumers[Idx].Queue.Count = 0 then dec(FCount); if FCount = 0 then begin - SetLength(Arr, Length(FReceivers)); + SetLength(Arr, Length(FConsumers)); Arr[Idx] := Value; - FCount := Length(FReceivers); - for var i := 0 to High(FReceivers) do + FCount := Length(FConsumers); + for var i := 0 to High(FConsumers) do if i <> Idx then begin - Arr[i] := FReceivers[i].Queue.Dequeue; - if FReceivers[i].Queue.Count > 0 then + Arr[i] := FConsumers[i].Queue.Dequeue; + if FConsumers[i].Queue.Count > 0 then dec(FCount); end; end else - FReceivers[Idx].Queue.Enqueue(Value); + FConsumers[Idx].Queue.Enqueue(Value); finally FLock.Exit; end; if Arr <> nil then - FDataProvider.Broadcast(Arr); + FContainedProvider.Broadcast(Arr); end; +{ TMycGenericAggregator } + constructor TMycGenericAggregator.Create(const AFunc: TConstFunc.TBroadcastProc, TState>); begin inherited Create; FFunc := AFunc; end; -function TMycGenericAggregator.ProcessData(const Value: S): TState; +function TMycGenericAggregator.Consume(const Value: S): TState; begin - Result := FFunc(Value, function(const Value: T): TState begin Result := FDataProvider.Broadcast(Value); end); + Result := FFunc(Value, Broadcast); end; -constructor TMycComposedConverter.Create(const AProcessor: IMycProcessor; const ADataProvider: TDataProvider); +{ TMycComposedConverter } + +constructor TMycComposedConverter.Create(const AConsumer: IConsumer; const AProducer: TProducer); begin inherited Create; - FProcessor := AProcessor; - FDataProvider := ADataProvider; + FConsumer := AConsumer; + FProducer := TProducer(AProducer); end; -function TMycComposedConverter.GetDataProvider: IDataProvider; +function TMycComposedConverter.GetConsumer: IConsumer; begin - Result := FDataProvider; + Result := FConsumer; end; -function TMycComposedConverter.ProcessData(const Value: S): TState; +function TMycComposedConverter.Link(const Consumer: IConsumer): TTag; begin - Result := FProcessor.ProcessData(Value); + // Delegate to the contained producer + Result := IProducer(FProducer).Link(Consumer); +end; + +procedure TMycComposedConverter.Unlink(Tag: TTag); +begin + // Delegate to the contained producer + IProducer(FProducer).Unlink(Tag); +end; + +{ TMycProducer } + +constructor TMycProducer.Create; +begin + inherited Create; +end; + +class constructor TMycProducer.CreateClass; +begin + FNull := TNull.Create; +end; + +destructor TMycProducer.Destroy; +begin + FListeners.Finalize; + inherited Destroy; +end; + +function TMycProducer.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.Consume(Value).Signal.Subscribe(done); + item := item.Prev; + end; + + Result := done.State; + finally + FListeners.Release; + end; +end; + +function TMycProducer.Link(const Consumer: IConsumer): TTag; +begin + // Add the Consumer to the notification list + FListeners.Lock; + try + Result := FListeners.Advise(Consumer); + finally + FListeners.Release; + end; +end; + +procedure TMycProducer.Unlink(Tag: TTag); +begin + FListeners.Lock; + try + FListeners.Unadvise(Tag); + finally + FListeners.Release; + end; +end; + +function TMycConsumer.TNull.Consume(const Value: T): TState; +begin + Result := TState.Null; +end; + +function TMycProducer.TNull.Link(const Consumer: IConsumer): TTag; +begin + Result := nil; +end; + +procedure TMycProducer.TNull.Unlink(Tag: TTag); +begin + +end; + +constructor TMycConverter.TConsumer.Create(AOwner: TMycConverter); +begin + inherited Create(AOwner); + FOwner := AOwner; +end; + +function TMycConverter.TConsumer.Consume(const Value: S): TState; +begin + Result := FOwner.Consume(Value); end; end. diff --git a/Src/Myc.Trade.DataPoint.pas b/Src/Myc.Trade.DataPoint.pas index 23f4a82..91358ac 100644 --- a/Src/Myc.Trade.DataPoint.pas +++ b/Src/Myc.Trade.DataPoint.pas @@ -10,74 +10,68 @@ uses Myc.DataRecord; type - // A generic interface for components that process data of type T. - IMycProcessor = interface - function ProcessData(const Value: T): TState; + TTag = Pointer; + + // A generic interface for components that consume data of type T. + IConsumer = interface + function Consume(const Value: T): TState; end; - TTag = Pointer; - IDataProvider = interface - function Link(const Receiver: IMycProcessor): TTag; + IProducer = interface + function Link(const Consumer: IConsumer): TTag; procedure Unlink(Tag: TTag); end; - IConverter = interface(IMycProcessor) + // A converter is a producer that internally uses a consumer to transform data. + IConverter = interface(IProducer) {$region 'private'} - function GetDataProvider: IDataProvider; + function GetConsumer: IConsumer; {$endregion} - property DataProvider: IDataProvider read GetDataProvider; + property Consumer: IConsumer read GetConsumer; end; - // Interface helper for IDataProvider providing the null object pattern. - TDataProvider = record + // Interface helper for IProducer providing the null object pattern. + TProducer = record + public type TSubscription = record private - FDataProvider: IDataProvider; + FProducer: IProducer; FTag: TTag; public procedure Unlink; - property DataProvider: IDataProvider read FDataProvider; + property Producer: IProducer read FProducer; end; - - strict private - class var - FNull: IDataProvider; - class constructor CreateClass; - private - FDataProvider: IDataProvider; + FProducer: IProducer; + class function GetNull: IProducer; static; inline; public - constructor Create(const ADataProvider: IDataProvider); + constructor Create(const AProducer: IProducer); // Managed record operators - class operator Initialize(out Dest: TDataProvider); - class operator Implicit(const A: IDataProvider): TDataProvider; overload; - class operator Implicit(const A: TDataProvider): IDataProvider; overload; + class operator Initialize(out Dest: TProducer); + class operator Implicit(const A: IProducer): TProducer; overload; + class operator Implicit(const A: TProducer): IProducer; overload; - function Link(const Receiver: IMycProcessor): TSubscription; inline; + function Link(const Consumer: IConsumer): TSubscription; inline; - function Chain(const Next: IMycProcessor): IMycProcessor; overload; inline; - function Chain(const Next: IConverter): TDataProvider; overload; inline; - function Chain(const Func: TConstFunc): TDataProvider; overload; inline; + function Chain(const Next: IConsumer): IConsumer; overload; inline; + + function Chain(const Next: IConverter): TProducer; overload; inline; + function Chain(const Func: TConstFunc): TProducer; overload; inline; function CreateEndpoint(Lookback: Int64): TLazy>; - // Extracts the field of a record by it's name (using RTTI). - function Field(const FieldName: String): TDataProvider; inline; + function CreateSequence(Count: Integer): TArray>; overload; experimental; - function MakeParallel: TDataProvider; inline; + // Extracts the field of a record by it's name (using RTTI). + function Field(const FieldName: String): TProducer; inline; + + function MakeParallel: TProducer; inline; // Provides access to the null object instance. - class property Null: IDataProvider read FNull; - end; - - IMycDataSequence = interface(IMycProcessor) - function GetCount: Integer; - function GetDataProvider(Idx: Integer): TDataProvider; - property Count: Integer read GetCount; - property DataProvider[Idx: Integer]: TDataProvider read GetDataProvider; default; + class property Null: IProducer read GetNull; end; // Interface helper for IConverter providing the null object pattern. @@ -85,14 +79,11 @@ type public type TBroadcastProc = reference to function(const Value: T): TState; - - strict private - class var - FNull: IConverter; - class constructor CreateClass; private FConverter: IConverter; - function GetDataProvider: TDataProvider; inline; + function GetConsumer: IConsumer; inline; + function GetProducer: TProducer; inline; + class function GetNull: IConverter; static; inline; public constructor Create(const AConverter: IConverter); @@ -101,17 +92,14 @@ type class operator Implicit(const A: IConverter): TConverter; overload; class operator Implicit(const A: TConverter): IConverter; overload; - class function Construct(const Processor: IMycProcessor; const DataProvider: TDataProvider): TConverter; static; + class function Construct(const Consumer: IConsumer; const Producer: TProducer): TConverter; static; class function CreateGeneric(const Func: TConstFunc): TConverter; static; class function CreateAggregation(const Func: TConstFunc): TConverter; static; - function CreateSequence(Count: Integer): IMycDataSequence; overload; - function Sequence(const Items: TArray>): IMycDataSequence; overload; + property Consumer: IConsumer read GetConsumer; + property Producer: TProducer read GetProducer; - // Provides access to the null object instance. - class property Null: IConverter read FNull; - // Wrapper for IConverter.DataProvider - property DataProvider: TDataProvider read GetDataProvider; + class property Null: IConverter read GetNull; end; // Factory for creating specific converter instances. @@ -126,12 +114,10 @@ type class function CreateDataPointConverter(const Func: TConstFunc): TConverter, TDataPoint>; static; - class function CreateSequence(Count: Integer; const Parent: TDataProvider): TArray>; overload; static; - class function FieldToRecord(const Layout: TDataRecord.TLayout; const Name: String): TConverter; static; class function FieldOfRecord(const Layout: TDataRecord.TLayout; const Name: String): TConverter; static; - class function Join(const DataProviders: TArray>): TDataProvider>; static; + class function Join(const Producers: TArray>): TProducer>; static; type TRecordMapping = record @@ -149,7 +135,7 @@ type class function JoinRecords( const TargetLayout: TDataRecord.TLayout; const Mapping: TArray; - const DataProviders: TArray> + const Producers: TArray> ): TConverter, TDataRecord>; static; end; @@ -158,100 +144,98 @@ implementation uses Myc.Trade.DataPoint.Impl; -{ TDataProvider } +{ TProducer } -procedure TDataProvider.TSubscription.Unlink; +procedure TProducer.TSubscription.Unlink; begin if FTag <> nil then begin - FDataProvider.Unlink(FTag); + FProducer.Unlink(FTag); FTag := nil; end; end; -class constructor TDataProvider.CreateClass; +constructor TProducer.Create(const AProducer: IProducer); begin - // Create the singleton null object instance. - FNull := TNullDataProvider.Create; + FProducer := AProducer; + if not Assigned(FProducer) then + FProducer := Null; end; -constructor TDataProvider.Create(const ADataProvider: IDataProvider); +function TProducer.Chain(const Next: IConsumer): IConsumer; begin - FDataProvider := ADataProvider; - if not Assigned(FDataProvider) then - FDataProvider := FNull; -end; - -function TDataProvider.Chain(const Next: IMycProcessor): IMycProcessor; -begin - FDataProvider.Link(Next); + FProducer.Link(Next); Result := Next; end; -function TDataProvider.Chain(const Next: IConverter): TDataProvider; +function TProducer.Chain(const Next: IConverter): TProducer; begin - FDataProvider.Link(Next); - Result := Next.DataProvider; + // REFACTOR: Link to the consumer part, return the converter as the new producer + FProducer.Link(Next.Consumer); + Result := Next; end; -function TDataProvider.Chain(const Func: TConstFunc): TDataProvider; +function TProducer.Chain(const Func: TConstFunc): TProducer; begin Result := Chain(TMycGenericConverter.Create(Func)); end; -function TDataProvider.CreateEndpoint(Lookback: Int64): TLazy>; +function TProducer.CreateEndpoint(Lookback: Int64): TLazy>; begin - Result := TMycDataEndpoint.Create(FDataProvider, Lookback); + Result := TMycDataEndpoint.Create(FProducer, Lookback); end; -function TDataProvider.Field(const FieldName: String): TDataProvider; +function TProducer.CreateSequence(Count: Integer): TArray>; +begin + FProducer.Link(TMycSequence.CreateSequence(Count, Result)); +end; + +function TProducer.Field(const FieldName: String): TProducer; begin Result := Chain(TMycRecordFieldReader.Create(FieldName)); end; -class operator TDataProvider.Initialize(out Dest: TDataProvider); +class function TProducer.GetNull: IProducer; begin - Dest.FDataProvider := FNull; + Result := TMycProducer.Null; end; -class operator TDataProvider.Implicit(const A: IDataProvider): TDataProvider; +class operator TProducer.Initialize(out Dest: TProducer); +begin + Dest.FProducer := Null; +end; + +class operator TProducer.Implicit(const A: IProducer): TProducer; begin Result.Create(A); end; -class operator TDataProvider.Implicit(const A: TDataProvider): IDataProvider; +class operator TProducer.Implicit(const A: TProducer): IProducer; begin - Result := A.FDataProvider; + Result := A.FProducer; end; -function TDataProvider.Link(const Receiver: IMycProcessor): TSubscription; +function TProducer.Link(const Consumer: IConsumer): TSubscription; begin - Result.FDataProvider := FDataProvider; - Result.FTag := FDataProvider.Link(Receiver); + Result.FProducer := FProducer; + Result.FTag := FProducer.Link(Consumer); end; -function TDataProvider.MakeParallel: TDataProvider; +function TProducer.MakeParallel: TProducer; begin Result := Chain(TMycParallelConverter.Create as IConverter); end; -{ TConverter } - -class constructor TConverter.CreateClass; -begin - FNull := TNullConverter.Create; -end; - constructor TConverter.Create(const AConverter: IConverter); begin FConverter := AConverter; if not Assigned(FConverter) then - FConverter := FNull; + FConverter := Null; end; -class function TConverter.Construct(const Processor: IMycProcessor; const DataProvider: TDataProvider): TConverter; +class function TConverter.Construct(const Consumer: IConsumer; const Producer: TProducer): TConverter; begin - Result := TMycComposedConverter.Create(Processor, DataProvider); + Result := TMycComposedConverter.Create(Consumer, Producer); end; class function TConverter.CreateAggregation(const Func: TConstFunc): TConverter; @@ -264,30 +248,24 @@ begin Result := TMycGenericConverter.Create(Func); end; -function TConverter.GetDataProvider: TDataProvider; +function TConverter.GetConsumer: IConsumer; begin - Result := FConverter.DataProvider; + Result := FConverter.Consumer; end; -function TConverter.CreateSequence(Count: Integer): IMycDataSequence; +function TConverter.GetProducer: TProducer; begin - Result := TMycSequence.Create(Count); - FConverter.DataProvider.Link(Result); + Result := FConverter; end; -function TConverter.Sequence(const Items: TArray>): IMycDataSequence; +class function TConverter.GetNull: IConverter; begin - var seq: IMycDataSequence := TMycSequence.Create(Length(Items)); - - for var i := 0 to High(Items) do - seq[i].Link(Items[i]); - - Result := seq; + Result := TMycConverter.Null; end; class operator TConverter.Initialize(out Dest: TConverter); begin - Dest.FConverter := FNull; + Dest.FConverter := Null; end; class operator TConverter.Implicit(const A: IConverter): TConverter; @@ -340,19 +318,6 @@ begin Result := TMycTicker.Create; end; -class function TConverter.CreateSequence(Count: Integer; const Parent: TDataProvider): TArray>; -begin - var seq: IMycDataSequence := TMycSequence.Create(Count); - - SetLength(Result, Count); - for var i := 0 to High(Result) do - begin - Result[i] := TMycIdentityConverter.Create; - seq[i].Link(Result[i]); - end; - Parent.Link(seq); -end; - class function TConverter.CreateOhlcAggregation(Timeframe: TTimeframe): TConverter, TDataPoint>; begin Result := TOhlcAggregation.Create(Timeframe); @@ -409,22 +374,23 @@ begin ); end; -class function TConverter.Join(const DataProviders: TArray>): TDataProvider>; +class function TConverter.Join(const Producers: TArray>): TProducer>; begin - Result := TMycDataJoin.Create(DataProviders); + Result := TMycDataJoin.Create(Producers); end; class function TConverter.JoinRecords( const TargetLayout: TDataRecord.TLayout; const Mapping: TArray; - const DataProviders: TArray> + const Producers: TArray> ): TConverter, TDataRecord>; begin - var RecProvider := Join(DataProviders); + var RecProvider := Join(Producers); Result := DataMapping(Mapping, TargetLayout); - RecProvider.Link(Result); + // REFACTOR: Link to the consumer part of the mapping converter. + RecProvider.Link(Result.Consumer); end; end. diff --git a/Src/Myc.Trade.DataStream.pas b/Src/Myc.Trade.DataStream.pas index 291f9d0..0ed1661 100644 --- a/Src/Myc.Trade.DataStream.pas +++ b/Src/Myc.Trade.DataStream.pas @@ -20,7 +20,7 @@ type IDataServer = interface ['{1F8E5A9D-E92A-44C1-9F3F-C4B82A6E94B3}'] procedure ClearCache; - function ProcessData(const Symbol: String; const Terminated: TState; const Processor: IMycProcessor>>): TState; + function ProcessData(const Symbol: String; const Terminated: TState; const Processor: IConsumer>>): TState; function EnumerateSymbols: TFuture>; end; @@ -60,14 +60,14 @@ type function ProcessChunks( const DataChunks: TArray>>; Terminated: TState; - Processor: IMycProcessor>> + Processor: IConsumer>> ): TState; function ProcessFile( FileInfo: TAuraDataFile; const DataFile: TFuture>>>; const Terminated: TState; - Processor: IMycProcessor>> + Processor: IConsumer>> ): TState; strict private @@ -94,7 +94,7 @@ type // Load data file and split content into chunks. function LoadDataFile(const DataFile: TAuraDataFile): TFuture>>>; - function ProcessData(const Symbol: String; const Terminated: TState; const Processor: IMycProcessor>>): TState; + function ProcessData(const Symbol: String; const Terminated: TState; const Processor: IConsumer>>): TState; property Path: String read GetPath; end; @@ -338,7 +338,7 @@ end; function TAuraDataServer.ProcessData( const Symbol: String; const Terminated: TState; - const Processor: IMycProcessor>> + const Processor: IConsumer>> ): TState; begin var cProc := Processor; @@ -364,7 +364,7 @@ end; function TAuraDataServer.ProcessChunks( const DataChunks: TArray>>; Terminated: TState; - Processor: IMycProcessor>> + Processor: IConsumer>> ): TState; begin var callProcess := @@ -375,7 +375,7 @@ begin function: TState begin if not Terminated.IsSet then - Result := Processor.ProcessData(cChunk); + Result := Processor.Consume(cChunk); end; end; @@ -391,7 +391,7 @@ function TAuraDataServer.ProcessFile( FileInfo: TAuraDataFile; const DataFile: TFuture>>>; const Terminated: TState; - Processor: IMycProcessor>> + Processor: IConsumer>> ): TState; begin if not FileInfo.IsValid or Terminated.IsSet then