diff --git a/AuraTrader/MainForm.pas b/AuraTrader/MainForm.pas index f9f870e..a532e8a 100644 --- a/AuraTrader/MainForm.pas +++ b/AuraTrader/MainForm.pas @@ -111,7 +111,6 @@ type procedure NewWorkspace; function CurrLayout: T; procedure AlignControl(Control: TControl); - function CreateStrategy1(Timeframe: TTimeframe): IMycProcessor>; function CreateStrategy2(Timeframe: TTimeframe): IMycProcessor>; published @@ -286,221 +285,6 @@ begin Control.Align := TAlignLayout.Top; end; -function TForm1.CreateStrategy1(Timeframe: TTimeframe): IMycProcessor>; -type - TSignal = record - Sig: Double; - SL: Double; - Entry: Double; - pnl: Double; - end; -var - panel: TMycChart.TPanel; -begin - var ticker := TConverter.CreateIdentity>; - Result := ticker; - - var OhlcPoint := ticker.Chain>(TConverter.CreateOhlcAggregation(Timeframe)); - - var Ohlc := TConverter.CreateSequence(2, OhlcPoint.Field('Data').Sender); - - var Closes := Ohlc[0].Field('Close'); - - var Hull := Closes.Chain(TIndicators.CreateHMA(250)).MakeParallel; - var Sma := Closes.Chain(TIndicators.CreateSMA(200)).MakeParallel; - - var Lowest: Double := Double.MaxValue; - var Highest: Double := Double.MinValue; - - var ATR := Ohlc[0].Chain(TIndicators.CreateATR(50)).MakeParallel; - - // next stage - - var ATREndPoint := TConverter.CreateEndpoint(ATR.Sender, 5); - var ATRSeries: TSeries; - - var HullEndPoint := TConverter.CreateEndpoint(Hull.Sender, 5); - var HullSeries: TSeries; - - var SmaEndPoint := TConverter.CreateEndpoint(Sma.Sender, 5); - var SmaSeries: TSeries; - - var curr: TSignal; - curr.SL := Double.NaN; - curr.Entry := Double.NaN; - - var Signal := - Ohlc[1] - .Chain( - function(const Ohlc: TOhlcItem): TSignal - begin - if Ohlc.Low < Lowest then - Lowest := Ohlc.Low; - if Ohlc.High > Highest then - Highest := Ohlc.High; - - Result := curr; - Result.Sig := 0; - var pnl: double := NaN; - - ATREndPoint.Update(ATRSeries); - HullEndPoint.Update(HullSeries); - SmaEndPoint.Update(SmaSeries); - - if (HullSeries[0] < SmaSeries[0]) and (HullSeries[1] >= SmaSeries[1]) then - begin - if curr.Sig > 0 then - pnl := Ohlc.Close - curr.Entry; - - curr.Sig := -1; - curr.SL := Highest; - curr.Entry := Ohlc.Close; - Result := curr; - end - else if (HullSeries[0] > SmaSeries[0]) and (HullSeries[1] <= SmaSeries[1]) then - begin - if curr.Sig < 0 then - pnl := curr.Entry - Ohlc.Close; - - curr.Sig := 1; - curr.SL := Lowest; - curr.Entry := Ohlc.Close; - Result := curr; - end; - - var atr := 15 * ATRSeries[0]; - if curr.Sig > 0 then - begin - if Ohlc.Close > curr.SL then - begin - if curr.SL < Ohlc.Close - atr then - curr.SL := Ohlc.Close - atr; - Result.SL := curr.SL; - end; - - if Ohlc.Low <= curr.SL then - begin - pnl := curr.SL - curr.Entry; - curr.Sig := 0; - Result.Sig := 0; - curr.SL := NaN; - end; - end - else if curr.Sig < 0 then - begin - if Ohlc.Close < curr.SL then - begin - if curr.SL > Ohlc.Close + atr then - curr.SL := Ohlc.Close + atr; - Result.SL := curr.SL; - end; - - if Ohlc.High >= curr.SL then - begin - pnl := curr.Entry - curr.SL; - curr.Sig := 0; - Result.Sig := 0; - curr.SL := NaN; - end; - end; - - if Result.Sig <> 0 then - begin - Lowest := Double.MaxValue; - Highest := Double.MinValue; - Result.SL := Double.NaN; - Result.Entry := Double.NaN; - end; - - Result.pnl := pnl; - end); - - var pnl := Signal.Field('pnl'); - - var FEquity: Double; - var FInit: Boolean; - - var equity := - TConverter.CreateAggregation( - function(const Value: Double; const Broadcast: TConverter.TBroadcastProc): TState - begin - if not FInit then - begin - FInit := true; - Broadcast(FEquity); - end; - - if not IsNan(Value) then - begin - FEquity := FEquity + Value; - Result := Broadcast(FEquity); - end; - end - ); - - // var equity: TConverter := TEquitySum.Create(10000); - pnl.Sender.Link(equity); - - var Layout := CurrLayout; - if Layout = nil then - exit; - - var Symbol := SelectedSymbol; - if Symbol = '' then - exit; - - var chart := TMycChart.Create(Self); - AlignControl(chart); - chart.Height := Layout.ChildrenRect.Width * 9 / 16; - chart.Lookback.Value := 50000; - - chart.SetXAxisSeries(M15, OhlcPoint.Field('Time').Sender); - - panel := chart.AddPanel; - panel.AddOhlcSeries(Ohlc[0].Sender); - panel.AddDoubleSeries(Hull.Sender, TAlphaColors.Cornflowerblue, 2); - panel.AddDoubleSeries(Sma.Sender, TAlphaColors.Brown, 1.5); - panel.AddDoubleSeries(Signal.Field('Entry').Sender, TAlphaColors.Green, 1); - panel.AddDoubleSeries(Signal.Field('SL').Sender, TAlphaColors.Red, 2); - - var mean := TConverter, Double>.CreateGeneric(TIndicators.CreateMean()); - - TConverter.Join([Hull.Sender, Sma.Sender]).Link(mean); - - panel.AddDoubleSeries(mean.Sender, TAlphaColors.Blue, 5); - - var pnlChart := TMycChart.Create(Self); - AlignControl(pnlChart); - pnlChart.Height := Layout.ChildrenRect.Width * 9 / 24; - pnlChart.Lookback.Value := 50000; - - pnlChart.SetXAxisCounter(equity.Sender); - - //////////// - var EMAFactory := TEMA.Create; - - var Params := TDataRecord.Create(EMAFactory.Params); - Params.SetValue('Period', 20); - - var indi := EMAFActory.CreateIndicator(Params); - - var EMAConv := TConverter.CreateGeneric(indi); - - var equityEMA := - equity - .Chain(TConverter.FieldToRecord(EMAFactory.Input, 'Price')) - .Chain(EMAConv) - .Chain(TConverter.FieldOfRecord(EMAFactory.Output, 'MA')); - - ////////////// - - panel := pnlChart.AddPanel; - panel.AddDoubleSeries(equity.Sender, TAlphaColors.Blue, 3); - panel.AddDoubleSeries(equityEMA.Sender, TAlphaColors.Gray, 2); - - ///// -end; - function TForm1.CreateStrategy2(Timeframe: TTimeframe): IMycProcessor>; type TSignal = record @@ -515,7 +299,7 @@ begin var ticker := TConverter.CreateIdentity>; Result := ticker; - var OhlcPoint := ticker.Chain>(TConverter.CreateOhlcAggregation(Timeframe)); + var OhlcPoint := ticker.Sender.Chain>(TConverter.CreateOhlcAggregation(Timeframe)); var Ohlc := OhlcPoint.Field('Data'); @@ -531,13 +315,13 @@ begin // next stage - var ATREndPoint := TConverter.CreateEndpoint(ATR.Sender, 5); + var ATREndPoint := ATR.CreateEndpoint(5); var ATRSeries: TSeries; - var HullEndPoint := TConverter.CreateEndpoint(Hull.Sender, 5); + var HullEndPoint := Hull.CreateEndpoint(5); var HullSeries: TSeries; - var SmaEndPoint := TConverter.CreateEndpoint(Sma.Sender, 5); + var SmaEndPoint := Sma.CreateEndpoint(5); var SmaSeries: TSeries; var curr: TSignal; @@ -546,10 +330,7 @@ begin var lastHull, lastSma: Double; - var conv := - TConverter.Join( - [Ohlc.Field('Low').Sender, Ohlc.Field('High').Sender, Closes.Sender, ATR.Sender, Hull.Sender, Sma.Sender] - ); + var conv := TConverter.Join([Ohlc.Field('Low'), Ohlc.Field('High'), Closes, ATR, Hull, Sma]); var Signal := TConverter, TSignal>.CreateGeneric( @@ -643,9 +424,9 @@ begin end ); - conv.Link(Signal); + conv.Chain(Signal); - var pnl := Signal.Field('pnl'); + var pnl := Signal.Sender.Field('pnl'); var FEquity: Double := 10000; var FInit: Boolean := false; @@ -668,7 +449,7 @@ begin end ); - pnl.Sender.Link(equity); + pnl.Chain(equity); var Layout := CurrLayout; if Layout = nil then @@ -683,18 +464,18 @@ begin chart.Height := Layout.ChildrenRect.Width * 9 / 16; chart.Lookback.Value := 50000; - chart.SetXAxisSeries(M15, OhlcPoint.Field('Time').Sender); + chart.SetXAxisSeries(M15, OhlcPoint.Field('Time')); panel := chart.AddPanel; - panel.AddOhlcSeries(Ohlc.Sender); - panel.AddDoubleSeries(Hull.Sender, TAlphaColors.Cornflowerblue, 2); - panel.AddDoubleSeries(Sma.Sender, TAlphaColors.Brown, 1.5); - panel.AddDoubleSeries(Signal.Field('Entry').Sender, TAlphaColors.Green, 1); - panel.AddDoubleSeries(Signal.Field('SL').Sender, TAlphaColors.Red, 2); + panel.AddOhlcSeries(Ohlc); + panel.AddDoubleSeries(Hull, TAlphaColors.Cornflowerblue, 2); + panel.AddDoubleSeries(Sma, TAlphaColors.Brown, 1.5); + panel.AddDoubleSeries(Signal.Sender.Field('Entry'), TAlphaColors.Green, 1); + panel.AddDoubleSeries(Signal.Sender.Field('SL'), TAlphaColors.Red, 2); var mean := TConverter, Double>.CreateGeneric(TIndicators.CreateMean()); - TConverter.Join([Hull.Sender, Sma.Sender]).Link(mean); + TConverter.Join([Hull, Sma]).Chain(mean); panel.AddDoubleSeries(mean.Sender, TAlphaColors.Blue, 5); @@ -717,6 +498,7 @@ begin var equityEMA := equity + .Sender .Chain(TConverter.FieldToRecord(EMAFactory.Input, 'Price')) .Chain(EMAConv) .Chain(TConverter.FieldOfRecord(EMAFactory.Output, 'MA')); @@ -725,7 +507,7 @@ begin panel := pnlChart.AddPanel; panel.AddDoubleSeries(equity.Sender, TAlphaColors.Blue, 3); - panel.AddDoubleSeries(equityEMA.Sender, TAlphaColors.Gray, 2); + panel.AddDoubleSeries(equityEMA, TAlphaColors.Gray, 2); ///// end; @@ -791,7 +573,7 @@ begin var ticker := TConverter.CreateTicker>; - ticker.Sender.Link(Processor); + ticker.Sender.Chain(Processor); FProcessDone := FProcessDone + FServer.ProcessData(Symbol, terminated, ticker); {$endif} @@ -823,7 +605,9 @@ begin ///// - var OhlcPoint := TConverter.CreateIdentity>; + var OhlcTicker := TConverter.CreateIdentity>; + + var OhlcPoint := OhlcTicker.Sender; var timeframe := TTimeframe.H; @@ -839,28 +623,28 @@ begin var Macd := Closes.MakeParallel.Chain(TIndicators.CreateMACD(12, 26, 9)); var Stoch := Ohlc.MakeParallel.Chain(TIndicators.CreateStochastic(14, 3)); - chart.SetXAxisSeries(timeframe, Timestamps.Sender); + chart.SetXAxisSeries(timeframe, Timestamps); var Panel := chart.AddPanel; - Panel.AddOhlcSeries(Ohlc.Sender); - Panel.AddDoubleSeries(Hull.Sender, TAlphaColors.Aliceblue); - Panel.AddDoubleSeries(Sma.Sender, TAlphaColors.Yellow); - Panel.AddDoubleSeries(Ema.Sender, TAlphaColors.Aqua); - Panel.AddDoubleSeries(Boli.Field('UpperBand').Sender, TAlphaColors.Gray); - Panel.AddDoubleSeries(Boli.Field('MiddleBand').Sender, TAlphaColors.Darkgray, 1.0); - Panel.AddDoubleSeries(Boli.Field('LowerBand').Sender, TAlphaColors.Gray); + Panel.AddOhlcSeries(Ohlc); + Panel.AddDoubleSeries(Hull, TAlphaColors.Aliceblue); + Panel.AddDoubleSeries(Sma, TAlphaColors.Yellow); + Panel.AddDoubleSeries(Ema, TAlphaColors.Aqua); + Panel.AddDoubleSeries(Boli.Field('UpperBand'), TAlphaColors.Gray); + Panel.AddDoubleSeries(Boli.Field('MiddleBand'), TAlphaColors.Darkgray, 1.0); + Panel.AddDoubleSeries(Boli.Field('LowerBand'), TAlphaColors.Gray); Panel := chart.AddPanel; - Panel.AddDoubleSeries(Rsi.Sender, TAlphaColors.Fuchsia); + Panel.AddDoubleSeries(Rsi, TAlphaColors.Fuchsia); Panel := chart.AddPanel; - Panel.AddDoubleSeries(Macd.Field('MacdLine').Sender, TAlphaColors.Orange); - Panel.AddDoubleSeries(Macd.Field('SignalLine').Sender, TAlphaColors.Dodgerblue); - Panel.AddDoubleSeries(Macd.Field('Histogram').Sender, TAlphaColors.Lightgreen); + Panel.AddDoubleSeries(Macd.Field('MacdLine'), TAlphaColors.Orange); + Panel.AddDoubleSeries(Macd.Field('SignalLine'), TAlphaColors.Dodgerblue); + Panel.AddDoubleSeries(Macd.Field('Histogram'), TAlphaColors.Lightgreen); Panel := chart.AddPanel; - Panel.AddDoubleSeries(Stoch.Field('K').Sender, TAlphaColors.Green); - Panel.AddDoubleSeries(Stoch.Field('D').Sender, TAlphaColors.Red); + Panel.AddDoubleSeries(Stoch.Field('K'), TAlphaColors.Green); + Panel.AddDoubleSeries(Stoch.Field('D'), TAlphaColors.Red); ///// { @@ -887,7 +671,7 @@ begin } ///// - ExecuteStrategy(Symbol, timeframe, OhlcPoint); + ExecuteStrategy(Symbol, timeframe, OhlcTicker); end; procedure TForm1.Strat2ButtonClick(Sender: TObject); @@ -897,7 +681,7 @@ begin exit; var timeframe := TTimeframe.M15; - // ExecuteStrategy(Symbol, timeframe, CreateStrategy2(timeframe)); + //ExecuteStrategy(Symbol, timeframe, CreateStrategy2(timeframe)); var tstStrat := StrategyTest.CreateStrategy1(timeframe); ExecuteStrategy(Symbol, timeframe, tstStrat); diff --git a/AuraTrader/StrategyTest.pas b/AuraTrader/StrategyTest.pas index 0423d81..322330b 100644 --- a/AuraTrader/StrategyTest.pas +++ b/AuraTrader/StrategyTest.pas @@ -10,7 +10,7 @@ uses Myc.DataRecord, Myc.Trade.Indicators; -function CreateStrategy1(Timeframe: TTimeframe): TConverter, Double>; +function CreateStrategy1(Timeframe: TTimeframe): TConverter, Double>; overload; implementation @@ -30,7 +30,7 @@ type begin var ticker := TConverter.CreateIdentity>; - var OhlcPoint := ticker.Chain>(TConverter.CreateOhlcAggregation(Timeframe)); + var OhlcPoint := ticker.Sender.Chain>(TConverter.CreateOhlcAggregation(Timeframe)); var Ohlc := OhlcPoint.Field('Data'); @@ -40,10 +40,7 @@ begin var Sma := Closes.Chain(TIndicators.CreateSMA(200)).MakeParallel; var ATR := Ohlc.Chain(TIndicators.CreateATR(50)).MakeParallel; - var conv := - TConverter.Join( - [Ohlc.Field('Low').Sender, Ohlc.Field('High').Sender, Closes.Sender, ATR.Sender, Hull.Sender, Sma.Sender] - ); + var conv := TConverter.Join([Ohlc.Field('Low'), Ohlc.Field('High'), Closes, ATR, Hull, Sma]); // STAGE 1: Signal Generation. This converter is stateless regarding the trade itself. // It only detects the crossover event and prepares the data for the next stage. @@ -52,42 +49,44 @@ begin var lastHull, lastSma: Double; var signalGenerator := - TConverter, TSignalEvent>.CreateGeneric( - function(const Values: TArray): TSignalEvent - begin - Result.Low := Values[0]; - Result.High := Values[1]; - Result.Close := Values[2]; - Result.ATR := Values[3]; - var hull := Values[4]; - var sma := Values[5]; - - if Result.Low < Lowest then - Lowest := Result.Low; - if Result.High > Highest then - Highest := Result.High; - - Result.Signal := 0; - Result.InitialSL := Double.NaN; - - if (hull < sma) and (lastHull >= lastSma) then + conv.Chain( + TConverter, TSignalEvent>.CreateGeneric( + function(const Values: TArray): TSignalEvent begin - Result.Signal := -1; - Result.InitialSL := Highest; - Highest := Double.MinValue; // Reset for next trend - Lowest := Double.MaxValue; + Result.Low := Values[0]; + Result.High := Values[1]; + Result.Close := Values[2]; + Result.ATR := Values[3]; + var hull := Values[4]; + var sma := Values[5]; + + if Result.Low < Lowest then + Lowest := Result.Low; + if Result.High > Highest then + Highest := Result.High; + + Result.Signal := 0; + Result.InitialSL := Double.NaN; + + if (hull < sma) and (lastHull >= lastSma) then + begin + Result.Signal := -1; + Result.InitialSL := Highest; + Highest := Double.MinValue; // Reset for next trend + Lowest := Double.MaxValue; + end + else if (hull > sma) and (lastHull <= lastSma) then + begin + Result.Signal := 1; + Result.InitialSL := Lowest; + Highest := Double.MinValue; // Reset for next trend + Lowest := Double.MaxValue; + end; + + lastHull := hull; + lastSma := sma; end - else if (hull > sma) and (lastHull <= lastSma) then - begin - Result.Signal := 1; - Result.InitialSL := Lowest; - Highest := Double.MinValue; // Reset for next trend - Lowest := Double.MaxValue; - end; - - lastHull := hull; - lastSma := sma; - end + ) ); // STAGE 2: Position Management. This stateful converter manages the lifecycle @@ -98,95 +97,93 @@ begin var currEntry := Double.NaN; var positionManager := - TConverter.CreateAggregation( - function(const Value: TSignalEvent; const Broadcast: TConverter.TBroadcastProc): TState - var - pnl: Double; - begin - Result := TState.Null; - pnl := Double.NaN; - - // 1. Check for a new signal to open or reverse a position - if Value.Signal <> 0 then + signalGenerator.Chain( + TConverter.CreateAggregation( + function(const Value: TSignalEvent; const Broadcast: TConverter.TBroadcastProc): TState + var + pnl: Double; begin - // If a position is already open, close it first - if currSig > 0 then - pnl := Value.Close - currEntry - else if currSig < 0 then - pnl := currEntry - Value.Close; + Result := TState.Null; + pnl := Double.NaN; - // Open new position - currSig := Value.Signal; - currEntry := Value.Close; - currSL := Value.InitialSL; - end - // 2. If no new signal, manage the currently open position - else - begin - var atrValue := 15 * Value.ATR; - if currSig > 0 then // Manage long position + // 1. Check for a new signal to open or reverse a position + if Value.Signal <> 0 then begin - if Value.Close > currSL then - if currSL < Value.Close - atrValue then - currSL := Value.Close - atrValue; + // If a position is already open, close it first + if currSig > 0 then + pnl := Value.Close - currEntry + else if currSig < 0 then + pnl := currEntry - Value.Close; - if Value.Low <= currSL then - begin - pnl := currSL - currEntry; - currSig := 0; // Close position - end; + // Open new position + currSig := Value.Signal; + currEntry := Value.Close; + currSL := Value.InitialSL; end - else if currSig < 0 then // Manage short position + // 2. If no new signal, manage the currently open position + else begin - if Value.Close < currSL then - if currSL > Value.Close + atrValue then - currSL := Value.Close + atrValue; - - if Value.High >= currSL then + var atrValue := 15 * Value.ATR; + if currSig > 0 then // Manage long position begin - pnl := currEntry - currSL; - currSig := 0; // Close position + if Value.Close > currSL then + if currSL < Value.Close - atrValue then + currSL := Value.Close - atrValue; + + if Value.Low <= currSL then + begin + pnl := currSL - currEntry; + currSig := 0; // Close position + end; + end + else if currSig < 0 then // Manage short position + begin + if Value.Close < currSL then + if currSL > Value.Close + atrValue then + currSL := Value.Close + atrValue; + + if Value.High >= currSL then + begin + pnl := currEntry - currSL; + currSig := 0; // Close position + end; end; end; - end; - // 3. If a PnL was generated (trade closed), broadcast it - if not IsNan(pnl) then - begin - currSL := Double.NaN; - Broadcast(pnl); - end; - end + // 3. If a PnL was generated (trade closed), broadcast it + if not IsNan(pnl) then + begin + currSL := Double.NaN; + Broadcast(pnl); + end; + end + ) ); - // Chain the stages together - conv.Link(signalGenerator); - signalGenerator.Sender.Link(positionManager); - // The final equity calculation remains the same, it just consumes the PnL from the position manager var FEquity: Double := 10000; var FInit: Boolean := false; var equity := - TConverter.CreateAggregation( - function(const Value: Double; const Broadcast: TConverter.TBroadcastProc): TState - begin - if not FInit then + positionManager.Chain( + TConverter.CreateAggregation( + function(const Value: Double; const Broadcast: TConverter.TBroadcastProc): TState begin - FInit := true; - Broadcast(FEquity); - end; + if not FInit then + begin + FInit := true; + Broadcast(FEquity); + end; - if not IsNan(Value) then - begin - FEquity := FEquity + Value; - Result := Broadcast(FEquity); - end; - end + if not IsNan(Value) then + begin + FEquity := FEquity + Value; + Result := Broadcast(FEquity); + end; + end + ) ); - positionManager.Sender.Link(equity); - - Result := TConverter, Double>.Construct(ticker, equity.Sender); + Result := TConverter, Double>.Construct(ticker, equity); end; end. diff --git a/Src/Myc.Fmx.Chart.Series.pas b/Src/Myc.Fmx.Chart.Series.pas index f2991b9..4182e57 100644 --- a/Src/Myc.Fmx.Chart.Series.pas +++ b/Src/Myc.Fmx.Chart.Series.pas @@ -272,7 +272,7 @@ begin inherited Create(AParent); FDataProvider := ADataProvider; - FReceiver := TConverter.CreateEndpoint(FDataProvider, Owner.Lookback.Value); + FReceiver := FDataProvider.CreateEndpoint(Owner.Lookback.Value); end; function TChartCustomLayer.GetCount: Int64; @@ -296,7 +296,7 @@ constructor TChartXAxisLayer.Create(AOwner: TMycChart; const ADataProvider: T begin inherited Create(AOwner); FDataProvider := ADataProvider; - FReceiver := TConverter.CreateEndpoint(FDataProvider, Owner.Lookback.Value); + FReceiver := FDataProvider.CreateEndpoint(Owner.Lookback.Value); end; function TChartXAxisLayer.GetCaption(Idx: Int64): String; diff --git a/Src/Myc.Fmx.Chart.pas b/Src/Myc.Fmx.Chart.pas index 7c296d0..ab334cf 100644 --- a/Src/Myc.Fmx.Chart.pas +++ b/Src/Myc.Fmx.Chart.pas @@ -636,9 +636,7 @@ function TMycChart.SetXAxisCounter(const DataProvider: TDataProvider): TMy begin FXAxisSeries.Free; - var counter := TConverter.CreateCounter; - DataProvider.Link(counter); - FXAxisSeries := TChartXAxisLayer.Create(Self, counter.Sender); + FXAxisSeries := TChartXAxisLayer.Create(Self, DataProvider.Chain(TConverter.CreateCounter)); Result := FXAxisSeries; end; diff --git a/Src/Myc.Trade.DataPoint.Impl.pas b/Src/Myc.Trade.DataPoint.Impl.pas index 776ea49..01457cd 100644 --- a/Src/Myc.Trade.DataPoint.Impl.pas +++ b/Src/Myc.Trade.DataPoint.Impl.pas @@ -22,7 +22,7 @@ type end; // Concrete data provider that manages a list of processors (listeners). - TMycContainedDataProvider = class abstract(TContainedObject, TDataProvider.IDataProvider) + TMycContainedDataProvider = class abstract(TContainedObject, IDataProvider) private FListeners: TMycNotifyList>; public @@ -31,9 +31,9 @@ type // Notifies all linked processors. function Broadcast(const Value: T): TState; // Link a Processor - function Link(const Processor: IMycProcessor): TDataProvider.TTag; + function Link(const Processor: IMycProcessor): TTag; // Unlink a linked strategy - procedure Unlink(Tag: TDataProvider.TTag); + procedure Unlink(Tag: TTag); end; TMycSequence = class(TMycProcessor, IMycDataSequence) @@ -50,17 +50,17 @@ type end; // Null object implementation for IDataProvider. - TNullDataProvider = class(TInterfacedObject, TDataProvider.IDataProvider) + TNullDataProvider = class(TInterfacedObject, IDataProvider) public - function Link(const Receiver: IMycProcessor): TDataProvider.TTag; - procedure Unlink(Tag: TDataProvider.TTag); + 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, TConverter.IConverter) + TMycConverter = class abstract(TMycProcessor, IConverter) private FSender: TMycContainedDataProvider; - function GetSender: TDataProvider.IDataProvider; + function GetSender: IDataProvider; protected function ProcessData(const Value: S): TState; override; abstract; // Broadcasts the given data to all linked processors. @@ -68,13 +68,13 @@ type public constructor Create; destructor Destroy; override; - property Sender: TDataProvider.IDataProvider read GetSender; + property Sender: IDataProvider read GetSender; end; // Null object implementation for IConverter. - TNullConverter = class(TInterfacedObject, TConverter.IConverter) + TNullConverter = class(TInterfacedObject, IConverter) private - function GetSender: TDataProvider.IDataProvider; + function GetSender: IDataProvider; public function ProcessData(const Value: S): TState; end; @@ -174,8 +174,8 @@ type end; private FProcessor: TMycContainedProcessor; - FTag: TDataProvider.TTag; - FDataProvider: TDataProvider; + FTag: TTag; + FDataProvider: IDataProvider; FLookback: Int64; FChanged: TFlag; FLock: TLightweightMREW; @@ -225,12 +225,12 @@ type end; // Endpoint that collects data into a series. - TMycDataJoin = class(TInterfacedObject, TDataProvider>.IDataProvider) + TMycDataJoin = class(TInterfacedObject, IDataProvider>) private FReceivers: array of record - DataProvider: TDataProvider; + DataProvider: IDataProvider; Processor: TMycContainedProcessor; - Tag: TDataProvider.TTag; + Tag: TTag; Queue: TQueue; end; @@ -241,16 +241,16 @@ type public constructor Create(const ADataProviders: TArray>); destructor Destroy; override; - property Sender: TMycContainedDataProvider> read FSender implements TDataProvider>.IDataProvider; + property Sender: TMycContainedDataProvider> read FSender implements IDataProvider>; end; - TMycComposedConverter = class(TMycProcessor, TConverter.IConverter) + TMycComposedConverter = class(TMycProcessor, IConverter) private FProcessor: IMycProcessor; FDataProvider: TDataProvider; protected function ProcessData(const Value: S): TState; override; - function GetSender: TDataProvider.IDataProvider; + function GetSender: IDataProvider; public constructor Create(const AProcessor: IMycProcessor; const ADataProvider: TDataProvider); end; @@ -314,7 +314,7 @@ begin end; end; -function TMycContainedDataProvider.Link(const Processor: IMycProcessor): TDataProvider.TTag; +function TMycContainedDataProvider.Link(const Processor: IMycProcessor): TTag; begin // Add the Processor to the notification list FListeners.Lock; @@ -325,7 +325,7 @@ begin end; end; -procedure TMycContainedDataProvider.Unlink(Tag: TDataProvider.TTag); +procedure TMycContainedDataProvider.Unlink(Tag: TTag); begin FListeners.Lock; try @@ -337,12 +337,12 @@ end; { TNullDataProvider } -function TNullDataProvider.Link(const Receiver: IMycProcessor): TDataProvider.TTag; +function TNullDataProvider.Link(const Receiver: IMycProcessor): TTag; begin Result := nil; end; -procedure TNullDataProvider.Unlink(Tag: TDataProvider.TTag); +procedure TNullDataProvider.Unlink(Tag: TTag); begin // Do nothing in the null implementation. end; @@ -366,14 +366,14 @@ begin Result := FSender.Broadcast(Value); end; -function TMycConverter.GetSender: TDataProvider.IDataProvider; +function TMycConverter.GetSender: IDataProvider; begin Result := FSender; end; { TNullConverter } -function TNullConverter.GetSender: TDataProvider.IDataProvider; +function TNullConverter.GetSender: IDataProvider; begin Result := TDataProvider.Null; end; @@ -499,7 +499,6 @@ begin FProcessor := TMycContainedProcessor.Create(Self, ProcessData); FTag := FDataProvider.Link(FProcessor); - end; destructor TMycDataEndpoint.Destroy; @@ -930,7 +929,7 @@ begin FDataProvider := ADataProvider; end; -function TMycComposedConverter.GetSender: TDataProvider.IDataProvider; +function TMycComposedConverter.GetSender: IDataProvider; begin Result := FDataProvider; end; diff --git a/Src/Myc.Trade.DataPoint.pas b/Src/Myc.Trade.DataPoint.pas index 93bc30e..bf7ac97 100644 --- a/Src/Myc.Trade.DataPoint.pas +++ b/Src/Myc.Trade.DataPoint.pas @@ -15,38 +15,60 @@ type function ProcessData(const Value: T): TState; end; + TTag = Pointer; + IDataProvider = interface + function Link(const Receiver: IMycProcessor): TTag; + procedure Unlink(Tag: TTag); + end; + + IConverter = interface(IMycProcessor) + function GetSender: IDataProvider; + property Sender: IDataProvider read GetSender; + end; + // Interface helper for IDataProvider providing the null object pattern. TDataProvider = record - public type - TTag = Pointer; - IDataProvider = interface - function Link(const Receiver: IMycProcessor): TTag; - procedure Unlink(Tag: TTag); + TLink = record + private + FDataProvider: IDataProvider; + FTag: TTag; + public + procedure Unlink; + property DataProvider: IDataProvider read FDataProvider; end; strict private class var - FNull: IDataProvider; + FNull: IDataProvider; class constructor CreateClass; private - FDataProvider: IDataProvider; + FDataProvider: IDataProvider; public - constructor Create(const ADataProvider: IDataProvider); + constructor Create(const ADataProvider: IDataProvider); // 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 Implicit(const A: IDataProvider): TDataProvider; overload; + class operator Implicit(const A: TDataProvider): IDataProvider; overload; - // Wrapper for IMycDataProvider methods - function Link(const Receiver: IMycProcessor): TTag; inline; - procedure Unlink(Tag: TTag); inline; + function Link(const Receiver: IMycProcessor): TLink; 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 CreateEndpoint(Lookback: Int64): TLazy>; + + // Extracts the field of a record by it's name (using RTTI). + function Field(const FieldName: String): TDataProvider; inline; + + function MakeParallel: TDataProvider; inline; // Provides access to the null object instance. - class property Null: IDataProvider read FNull; + class property Null: IDataProvider read FNull; end; IMycDataSequence = interface(IMycProcessor) @@ -60,57 +82,38 @@ type TConverter = record public type - IConverter = interface(IMycProcessor) - {$region 'private'} - function GetSender: TDataProvider.IDataProvider; - {$endregion} - property Sender: TDataProvider.IDataProvider read GetSender; - end; - TBroadcastProc = reference to function(const Value: T): TState; strict private class var - FNull: IConverter; + FNull: IConverter; class constructor CreateClass; private - FConverter: IConverter; + FConverter: IConverter; function GetSender: TDataProvider; inline; public - constructor Create(const AConverter: IConverter); + constructor Create(const AConverter: IConverter); // Managed record operators class operator Initialize(out Dest: TConverter); - class operator Implicit(const A: IConverter): TConverter; overload; - class operator Implicit(const A: TConverter): IConverter; overload; + 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 CreateGeneric(const Func: TConstFunc): TConverter; static; class function CreateAggregation(const Func: TConstFunc): TConverter; static; - class function CreateParallel(const Func: TConstFunc): TConverter; static; - - function Chain(const Next: TConverter): TConverter; overload; inline; - function Chain(const Func: TConstFunc): TConverter; overload; inline; - function ChainParallel(const Func: TConstFunc): TConverter; overload; inline; - - function MakeParallel: TConverter; overload; inline; - - // Extracts the field of a record by it's name (using RTTI). - function Field(const FieldName: String): TConverter; overload; inline; function Sequence(Count: Integer): IMycDataSequence; overload; - function Sequence(const Items: TArray): IMycDataSequence; overload; + function Sequence(const Items: TArray>): IMycDataSequence; overload; // Provides access to the null object instance. - class property Null: IConverter read FNull; + class property Null: IConverter read FNull; // Wrapper for IConverter.Sender property Sender: TDataProvider read GetSender; end; // Factory for creating specific converter instances. TConverter = record - class function CreateEndpoint(const DataProvider: TDataProvider; Lookback: Int64): TLazy>; static; - class function CreateCounter: TConverter; static; class function CreateTicker: TConverter, T>; static; class function CreateRecordField(const FieldName: String): TConverter; static; @@ -123,8 +126,6 @@ type class function CreateSequence(Count: Integer; const Parent: TDataProvider): TArray>; overload; static; - class function Parallel(Parent: TDataProvider): TConverter; 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; @@ -163,36 +164,64 @@ begin FNull := TNullDataProvider.Create; end; -constructor TDataProvider.Create(const ADataProvider: IDataProvider); +constructor TDataProvider.Create(const ADataProvider: IDataProvider); begin FDataProvider := ADataProvider; if not Assigned(FDataProvider) then FDataProvider := FNull; end; +function TDataProvider.Chain(const Next: IMycProcessor): IMycProcessor; +begin + FDataProvider.Link(Next); + Result := Next; +end; + +function TDataProvider.Chain(const Next: IConverter): TDataProvider; +begin + FDataProvider.Link(Next); + Result := Next.Sender; +end; + +function TDataProvider.Chain(const Func: TConstFunc): TDataProvider; +begin + Result := Chain(TMycGenericConverter.Create(Func)); +end; + +function TDataProvider.CreateEndpoint(Lookback: Int64): TLazy>; +begin + Result := TMycDataEndpoint.Create(FDataProvider, Lookback); +end; + +function TDataProvider.Field(const FieldName: String): TDataProvider; +begin + Result := Chain(TMycRecordFieldReader.Create(FieldName)); +end; + class operator TDataProvider.Initialize(out Dest: TDataProvider); begin Dest.FDataProvider := FNull; end; -class operator TDataProvider.Implicit(const A: IDataProvider): TDataProvider; +class operator TDataProvider.Implicit(const A: IDataProvider): TDataProvider; begin Result.Create(A); end; -class operator TDataProvider.Implicit(const A: TDataProvider): IDataProvider; +class operator TDataProvider.Implicit(const A: TDataProvider): IDataProvider; begin Result := A.FDataProvider; end; -function TDataProvider.Link(const Receiver: IMycProcessor): TTag; +function TDataProvider.Link(const Receiver: IMycProcessor): TLink; begin - Result := FDataProvider.Link(Receiver); + Result.FDataProvider := FDataProvider; + Result.FTag := FDataProvider.Link(Receiver); end; -procedure TDataProvider.Unlink(Tag: TTag); +function TDataProvider.MakeParallel: TDataProvider; begin - FDataProvider.Unlink(Tag); + Result := Chain(TMycParallelConverter.Create as IConverter); end; { TConverter } @@ -202,29 +231,13 @@ begin FNull := TNullConverter.Create; end; -constructor TConverter.Create(const AConverter: IConverter); +constructor TConverter.Create(const AConverter: IConverter); begin FConverter := AConverter; if not Assigned(FConverter) then FConverter := FNull; end; -function TConverter.Chain(const Next: TConverter): TConverter; -begin - FConverter.Sender.Link(Next); - Result := Next; -end; - -function TConverter.Chain(const Func: TConstFunc): TConverter; -begin - Result := Chain(TMycGenericConverter.Create(Func)); -end; - -function TConverter.ChainParallel(const Func: TConstFunc): TConverter; -begin - Result := Chain(TMycGenericParallelConverter.Create(Func)); -end; - class function TConverter.Construct(const Processor: IMycProcessor; const DataProvider: TDataProvider): TConverter; begin Result := TMycComposedConverter.Create(Processor, DataProvider); @@ -240,33 +253,18 @@ begin Result := TMycGenericConverter.Create(Func); end; -class function TConverter.CreateParallel(const Func: TConstFunc): TConverter; -begin - Result := TMycGenericParallelConverter.Create(Func); -end; - -function TConverter.Field(const FieldName: String): TConverter; -begin - Result := Chain(TConverter.CreateRecordField(FieldName)); -end; - function TConverter.GetSender: TDataProvider; begin Result := FConverter.Sender; end; -function TConverter.MakeParallel: TConverter; -begin - Result := Chain(TMycGenericParallelConverter.Create(function(const Value: T): T begin Result := Value; end)); -end; - function TConverter.Sequence(Count: Integer): IMycDataSequence; begin Result := TMycSequence.Create(Count); FConverter.Sender.Link(Result); end; -function TConverter.Sequence(const Items: TArray): IMycDataSequence; +function TConverter.Sequence(const Items: TArray>): IMycDataSequence; begin var seq: IMycDataSequence := TMycSequence.Create(Length(Items)); @@ -281,12 +279,12 @@ begin Dest.FConverter := FNull; end; -class operator TConverter.Implicit(const A: IConverter): TConverter; +class operator TConverter.Implicit(const A: IConverter): TConverter; begin Result.Create(A); end; -class operator TConverter.Implicit(const A: TConverter): IConverter; +class operator TConverter.Implicit(const A: TConverter): IConverter; begin Result := A.FConverter; end; @@ -316,11 +314,6 @@ begin ); end; -class function TConverter.CreateEndpoint(const DataProvider: TDataProvider; Lookback: Int64): TLazy>; -begin - Result := TMycDataEndpoint.Create(DataProvider, Lookback); -end; - class function TConverter.CreateIdentity: TConverter; begin Result := TMycIdentityConverter.Create; @@ -423,10 +416,13 @@ begin RecProvider.Link(Result); end; -class function TConverter.Parallel(Parent: TDataProvider): TConverter; +procedure TDataProvider.TLink.Unlink; begin - Result := TMycParallelConverter.Create; - Parent.Link(Result); + if FTag <> nil then + begin + FDataProvider.Unlink(FTag); + FTag := nil; + end; end; end.