From aa53a8895307ff8d62ae03a5cec60bd36f8fb1d0 Mon Sep 17 00:00:00 2001 From: Michael Schimmel Date: Fri, 25 Jul 2025 11:54:53 +0200 Subject: [PATCH] Unit refactoring Fixed massive heap corruption bug in TDataRecord --- AuraTrader/AuraTrader.dpr | 3 +- AuraTrader/AuraTrader.dproj | 1 + AuraTrader/MainForm.pas | 7 +- AuraTrader/StrategyTest.pas | 7 +- Src/Myc.Data.Pipeline.Impl.pas | 267 ++++---------------------------- Src/Myc.Data.Pipeline.pas | 35 ++--- Src/Myc.Data.Records.pas | 115 +++++++++----- Src/Myc.Trade.Indicators.pas | 87 ++++++----- Src/Myc.Trade.Pipeline.Impl.pas | 218 ++++++++++++++++++++++++++ Src/Myc.Trade.Pipeline.pas | 30 ++++ Src/Myc.Trade.Types.pas | 8 +- Test/TestDataArray.pas | 2 +- Test/TestDataRecord.pas | 40 ++++- 13 files changed, 461 insertions(+), 359 deletions(-) create mode 100644 Src/Myc.Trade.Pipeline.Impl.pas create mode 100644 Src/Myc.Trade.Pipeline.pas diff --git a/AuraTrader/AuraTrader.dpr b/AuraTrader/AuraTrader.dpr index 31edf55..19e7cf9 100644 --- a/AuraTrader/AuraTrader.dpr +++ b/AuraTrader/AuraTrader.dpr @@ -6,7 +6,8 @@ uses FMX.Forms, MainForm in 'MainForm.pas' {Form1}, TestModule in 'TestModule.pas', - DynamicFMXControl in 'DynamicFMXControl.pas'; + DynamicFMXControl in 'DynamicFMXControl.pas', + Myc.Trade.Pipeline.Impl in '..\Src\Myc.Trade.Pipeline.Impl.pas'; {$R *.res} diff --git a/AuraTrader/AuraTrader.dproj b/AuraTrader/AuraTrader.dproj index 3808279..15ce6c0 100644 --- a/AuraTrader/AuraTrader.dproj +++ b/AuraTrader/AuraTrader.dproj @@ -137,6 +137,7 @@ + Base diff --git a/AuraTrader/MainForm.pas b/AuraTrader/MainForm.pas index 0f6102e..9a33a0b 100644 --- a/AuraTrader/MainForm.pas +++ b/AuraTrader/MainForm.pas @@ -32,6 +32,7 @@ uses Myc.Futures, Myc.Trade.Types, Myc.Trade.DataStream, + Myc.Trade.Pipeline, Myc.Data.Series, Myc.Data.Pipeline, Myc.Signals, @@ -298,7 +299,7 @@ begin var ticker := TConverter.CreateIdentity>; Result := ticker.Consumer; - var OhlcPoint := ticker.Producer.Chain>(TConverter.CreateOhlcAggregation(Timeframe)); + var OhlcPoint := ticker.Producer.Chain>(TTradeConverter.CreateOhlcAggregation(Timeframe)); var Ohlc := OhlcPoint.Field('Data'); @@ -423,7 +424,7 @@ begin var equity := TConverter.CreateAggregation( - function(const Value: Double; const Broadcast: TConverter.TBroadcastProc): TState + function(const Value: Double; const Broadcast: TBroadcastFunc): TState begin if not FInit then begin @@ -599,7 +600,7 @@ begin var ticker := TConverter.CreateIdentity>; - var OhlcPoint := ticker.Producer.Chain>(TConverter.CreateOhlcAggregation(timeframe)); + var OhlcPoint := ticker.Producer.Chain>(TTradeConverter.CreateOhlcAggregation(timeframe)); // var OhlcTicker := TConverter.CreateIdentity>; // var OhlcPoint := OhlcTicker.Sender; diff --git a/AuraTrader/StrategyTest.pas b/AuraTrader/StrategyTest.pas index 0d5a999..14e72d9 100644 --- a/AuraTrader/StrategyTest.pas +++ b/AuraTrader/StrategyTest.pas @@ -6,6 +6,7 @@ uses Myc.Signals, Myc.Data.Pipeline, Myc.Trade.Types, + Myc.Trade.Pipeline, Myc.Trade.Indicators; function CreateStrategy1(Timeframe: TTimeframe): TConverter, Double>; overload; @@ -28,7 +29,7 @@ type begin var ticker := TConverter.CreateIdentity>; - var OhlcPoint := ticker.Producer.Chain>(TConverter.CreateOhlcAggregation(Timeframe)); + var OhlcPoint := ticker.Producer.Chain>(TTradeConverter.CreateOhlcAggregation(Timeframe)); var Ohlc := OhlcPoint.Field('Data'); @@ -97,7 +98,7 @@ begin var positionManager := signalGenerator.Chain( TConverter.CreateAggregation( - function(const Value: TSignalEvent; const Broadcast: TConverter.TBroadcastProc): TState + function(const Value: TSignalEvent; const Broadcast: TBroadcastFunc): TState var pnl: Double; begin @@ -164,7 +165,7 @@ begin var equity := positionManager.Chain( TConverter.CreateAggregation( - function(const Value: Double; const Broadcast: TConverter.TBroadcastProc): TState + function(const Value: Double; const Broadcast: TBroadcastFunc): TState begin if not FInit then begin diff --git a/Src/Myc.Data.Pipeline.Impl.pas b/Src/Myc.Data.Pipeline.Impl.pas index 0747a15..71cc27c 100644 --- a/Src/Myc.Data.Pipeline.Impl.pas +++ b/Src/Myc.Data.Pipeline.Impl.pas @@ -9,9 +9,8 @@ uses Myc.Signals, Myc.Mutable, Myc.Core.Notifier, - Myc.Data.Pipeline, - Myc.Trade.Types, - Myc.Data.Series; + Myc.Data.Series, + Myc.Data.Pipeline; type // Abstract base class for data consumers. @@ -31,16 +30,6 @@ type class property Null: IConsumer read FNull; end; - // A consumer implementation that is owned by a controller. - TMycGenericConsumer = class(TMycConsumer) - private - FProc: TConstFunc; - protected - function Consume(const Value: T): TState; override; final; - public - constructor Create(const Controller: IInterface; const AProc: TConstFunc); - end; - TMycProducer = class(TInterfacedObject, IProducer) strict private type @@ -63,16 +52,14 @@ type class property Null: IProducer read FNull; end; - // Concrete producer that manages a list of consumers (listeners). - TMycContainedProducer = class(TContainedObject, IProducer) + // A consumer implementation that is owned by a controller. + TMycGenericConsumer = class(TMycConsumer) private - FListeners: TMycNotifyList>; + FProc: TConvertFunc; + protected + function Consume(const Value: T): TState; override; final; public - constructor Create(const Controller: IInterface); - destructor Destroy; override; - function Broadcast(const Value: T): TState; - function Link(const Consumer: IConsumer): TTag; - procedure Unlink(Tag: TTag); + constructor Create(const Controller: IInterface; const AProc: TConvertFunc); end; // Abstract base class for components that now act as a producer and contain a consumer. @@ -101,33 +88,45 @@ type class property Null: IConverter read FNull; end; + // Concrete producer that manages a list of consumers (listeners). + TMycContainedProducer = class(TContainedObject, IProducer) + private + FListeners: TMycNotifyList>; + public + 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; + // A generic converter that uses a function reference for the conversion logic. TMycGenericConverter = class(TMycConverter) private - FFunc: TConstFunc; + FFunc: TConvertFunc; protected function Consume(const Value: S): TState; override; public - constructor Create(const AFunc: TConstFunc); + constructor Create(const AFunc: TConvertFunc); end; TMycGenericAggregator = class(TMycConverter) private - FFunc: TConstFunc.TBroadcastProc, TState>; + FFunc: TAggregateFunc; protected function Consume(const Value: S): TState; override; public - constructor Create(const AFunc: TConstFunc.TBroadcastProc, TState>); + constructor Create(const AFunc: TAggregateFunc); end; TMycGenericParallelConverter = class(TMycConverter) private - FFunc: TConstFunc; + FFunc: TConvertFunc; FQueue: TState; protected function Consume(const Value: S): TState; override; public - constructor Create(const AFunc: TConstFunc); + constructor Create(const AFunc: TConvertFunc); end; TMycIdentityConverter = class(TMycConverter) @@ -187,21 +186,6 @@ type class function CreateDataEndpoint(Lookback: Integer; out Series: TLazy>): IConsumer; end; - TTickAggregation = class(TMycConverter, TDataPoint>) - private - FTimeframe: TTimeframe; - FCurrentBar: TDataPoint; - 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); - property CurrentBar: TDataPoint read GetCurrentBar; - property Timeframe: TTimeframe read GetTimeframe; - end; - TMycParallelConverter = class(TMycConverter) private FQueue: TState; @@ -209,21 +193,6 @@ type function Consume(const Value: T): TState; override; final; end; - TOhlcAggregation = class(TMycConverter, TDataPoint>) - private - FTimeframe: TTimeframe; - FCurrentBar: TDataPoint; - 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); - property CurrentBar: TDataPoint read GetCurrentBar; - property Timeframe: TTimeframe read GetTimeframe; - end; - // Endpoint that collects data into a series. TMycDataJoin = class(TInterfacedObject, IProducer>) private @@ -273,11 +242,7 @@ type implementation uses - System.TypInfo, System.RTTI, - System.DateUtils, - System.Math, - Winapi.Windows, Myc.TaskManager; class constructor TMycConsumer.CreateClass; @@ -392,7 +357,7 @@ end; { TMycGenericConverter } -constructor TMycGenericConverter.Create(const AFunc: TConstFunc); +constructor TMycGenericConverter.Create(const AFunc: TConvertFunc); begin inherited Create; FFunc := AFunc; @@ -462,7 +427,7 @@ end; { TMycGenericConsumer } -constructor TMycGenericConsumer.Create(const Controller: IInterface; const AProc: TConstFunc); +constructor TMycGenericConsumer.Create(const Controller: IInterface; const AProc: TConvertFunc); begin inherited Create(Controller); FProc := AProc; @@ -550,94 +515,6 @@ begin end; end; -{ TTickAggregation } - -constructor TTickAggregation.Create(const ATimeframe: TTimeframe); -begin - inherited Create; - FTimeframe := ATimeframe; -end; - -function TTickAggregation.GetBarStartTime(const TimeStamp: TDateTime; const Timeframe: TTimeframe): TDateTime; -var - baseTime: TDateTime; -begin - // 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)); - M5: Result := RecodeMinute(RecodeSecond(baseTime, 0), MinuteOf(TimeStamp) - (MinuteOf(TimeStamp) mod 5)); - 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: Result := Floor(TimeStamp) - (Floor(TimeStamp) mod 2); - D3: Result := Floor(TimeStamp) - (Floor(TimeStamp) mod 3); - W: Result := TimeStamp.StartOfTheWeek; - MN: Result := TimeStamp.StartOfTheMonth; - MN3: Result := RecodeMonth(TimeStamp.StartOfTheMonth, (MonthOf(TimeStamp) - 1) div 3 * 3 + 1); - MN6: Result := RecodeMonth(TimeStamp.StartOfTheMonth, (MonthOf(TimeStamp) - 1) div 6 * 6 + 1); - Y: Result := TimeStamp.StartOfTheYear; - else - Result := 0; - end; -end; - -function TTickAggregation.GetCurrentBar: TDataPoint; -begin - Result := FCurrentBar; -end; - -function TTickAggregation.GetTimeframe: TTimeframe; -begin - Result := FTimeframe; -end; - -function TTickAggregation.Consume(const Value: TDataPoint): TState; -var - barStartTime: TDateTime; - lastBarTime: TDateTime; -begin - barStartTime := GetBarStartTime(Value.Time, FTimeframe); - lastBarTime := FCurrentBar.Time; - - if (barStartTime > lastBarTime) then - begin - if (lastBarTime > 0) then - begin - Result := Broadcast(FCurrentBar); - end; - - FCurrentBar.Data.Open := Value.Data; - FCurrentBar.Data.High := Value.Data; - FCurrentBar.Data.Low := Value.Data; - FCurrentBar.Data.Close := Value.Data; - FCurrentBar.Data.Volume := 1; - FCurrentBar.Time := barStartTime; - end - else - begin - 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; - FCurrentBar.Data.Volume := FCurrentBar.Data.Volume + 1; - end; -end; - { TMycSequence } constructor TMycSequence.Create(ACount: Integer); @@ -688,7 +565,7 @@ end; { TMycGenericParallelConverter } -constructor TMycGenericParallelConverter.Create(const AFunc: TConstFunc); +constructor TMycGenericParallelConverter.Create(const AFunc: TConvertFunc); begin inherited Create; FFunc := AFunc; @@ -711,88 +588,6 @@ begin FQueue := Result; end; -{ TOhlcAggregation } - -constructor TOhlcAggregation.Create(const ATimeframe: TTimeframe); -begin - inherited Create; - FTimeframe := ATimeframe; -end; - -function TOhlcAggregation.GetBarStartTime(const TimeStamp: TDateTime; const Timeframe: TTimeframe): TDateTime; -begin - // 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)); - M5: Result := RecodeMinute(RecodeSecond(baseTime, 0), MinuteOf(TimeStamp) - (MinuteOf(TimeStamp) mod 5)); - 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: Result := Floor(TimeStamp) - (Floor(TimeStamp) mod 2); - D3: Result := Floor(TimeStamp) - (Floor(TimeStamp) mod 3); - W: Result := TimeStamp.StartOfTheWeek; - MN: Result := TimeStamp.StartOfTheMonth; - MN3: Result := RecodeMonth(TimeStamp.StartOfTheMonth, (MonthOf(TimeStamp) - 1) div 3 * 3 + 1); - MN6: Result := RecodeMonth(TimeStamp.StartOfTheMonth, (MonthOf(TimeStamp) - 1) div 6 * 6 + 1); - Y: Result := TimeStamp.StartOfTheYear; - else - Result := 0; - end; -end; - -function TOhlcAggregation.GetCurrentBar: TDataPoint; -begin - Result := FCurrentBar; -end; - -function TOhlcAggregation.GetTimeframe: TTimeframe; -begin - Result := FTimeframe; -end; - -function TOhlcAggregation.Consume(const Value: TDataPoint): TState; -var - barStartTime: TDateTime; - lastBarTime: TDateTime; -begin - barStartTime := GetBarStartTime(Value.Time, FTimeframe); - lastBarTime := FCurrentBar.Time; - - if (barStartTime > lastBarTime) then - begin - if (lastBarTime > 0) then - begin - Result := Broadcast(FCurrentBar); - end; - - FCurrentBar.Data := Value.Data; - FCurrentBar.Time := barStartTime; - end - else - begin - 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; - FCurrentBar.Data.Volume := FCurrentBar.Data.Volume + Value.Data.Volume; - end; -end; - { TMycDataJoin } constructor TMycDataJoin.Create(ACount: Integer); @@ -807,7 +602,7 @@ begin FContainedProvider := TMycContainedProducer>.Create(Self); var cFunc := - function(Idx: Integer): TConstFunc + function(Idx: Integer): TConvertFunc begin Result := function(const Value: T): TState begin Result := Consume(Idx, Value); end end; @@ -882,7 +677,7 @@ end; { TMycGenericAggregator } -constructor TMycGenericAggregator.Create(const AFunc: TConstFunc.TBroadcastProc, TState>); +constructor TMycGenericAggregator.Create(const AFunc: TAggregateFunc); begin inherited Create; FFunc := AFunc; diff --git a/Src/Myc.Data.Pipeline.pas b/Src/Myc.Data.Pipeline.pas index 160133e..8e113d4 100644 --- a/Src/Myc.Data.Pipeline.pas +++ b/Src/Myc.Data.Pipeline.pas @@ -5,7 +5,6 @@ interface uses Myc.Signals, Myc.Mutable, - Myc.Trade.Types, Myc.Data.Series, Myc.Data.Records; @@ -31,6 +30,10 @@ type property Consumer: IConsumer read GetConsumer; end; + TConvertFunc = reference to function(const Value: S): T; + TBroadcastFunc = reference to function(const Value: T): TState; + TAggregateFunc = reference to function(const Value: S; const Broadcast: TBroadcastFunc): TState; + // Interface helper for IProducer providing the null object pattern nad subscriptions for linked consumers. TProducer = record public @@ -64,7 +67,7 @@ type // Chain consumers 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 Chain(const Func: TConvertFunc): TProducer; overload; inline; // Extracts the field of a record by it's name (using RTTI). function Field(const FieldName: String): TProducer; inline; @@ -76,9 +79,6 @@ type // Interface helper for IConverter providing the null object pattern. TConverter = record - public - type - TBroadcastProc = reference to function(const Value: T): TState; private FConverter: IConverter; function GetConsumer: IConsumer; inline; @@ -94,8 +94,8 @@ type class function Construct(const Consumer: IConsumer; const Producer: TProducer): TConverter; static; - class function CreateConverter(const Func: TConstFunc): TConverter; static; - class function CreateAggregation(const Func: TConstFunc): TConverter; static; + class function CreateConverter(const Func: TConvertFunc): TConverter; static; + class function CreateAggregation(const Func: TAggregateFunc): TConverter; static; property Consumer: IConsumer read GetConsumer; property Producer: TProducer read GetProducer; @@ -114,9 +114,6 @@ type class function CreateTicker: TConverter, T>; static; class function CreateIdentity: TConverter; static; - class function CreateTickAggregation(Timeframe: TTimeframe): TConverter, TDataPoint>; static; - class function CreateOhlcAggregation(Timeframe: TTimeframe): TConverter, TDataPoint>; static; - class function CreateEndpoint(Lookback: Int64; out Series: TLazy>): IConsumer; static; class function Join(const Producers: TArray>): TProducer>; static; @@ -188,7 +185,7 @@ begin Result := Next; end; -function TProducer.Chain(const Func: TConstFunc): TProducer; +function TProducer.Chain(const Func: TConvertFunc): TProducer; begin Result := Chain(TMycGenericConverter.Create(Func)); end; @@ -251,12 +248,12 @@ begin Result := TMycComposedConverter.Create(Consumer, Producer); end; -class function TConverter.CreateAggregation(const Func: TConstFunc): TConverter; +class function TConverter.CreateAggregation(const Func: TAggregateFunc): TConverter; begin Result := TMycGenericAggregator.Create(Func); end; -class function TConverter.CreateConverter(const Func: TConstFunc): TConverter; +class function TConverter.CreateConverter(const Func: TConvertFunc): TConverter; begin Result := TMycGenericConverter.Create(Func); end; @@ -291,11 +288,6 @@ begin Result := A.FConverter; end; -class function TConverter.CreateTickAggregation(Timeframe: TTimeframe): TConverter, TDataPoint>; -begin - Result := TTickAggregation.Create(Timeframe); -end; - { TConverter } class function TConverter.CreateCounter: TConverter; @@ -318,11 +310,6 @@ begin Result := TMycTicker.Create; end; -class function TConverter.CreateOhlcAggregation(Timeframe: TTimeframe): TConverter, TDataPoint>; -begin - Result := TOhlcAggregation.Create(Timeframe); -end; - class function TConverter.DataMapping( const Inputs: TArray; const Output: TDataRecord.TLayout @@ -350,7 +337,7 @@ begin Result := TDataRecord.Create(Output); for var i := 0 to High(Idxs) do for var j := 0 to High(Idxs[i]) do - Inputs[i].CopyField(Idxs[i][j].FromIdx, Result, Idxs[i][j].ToIdx); + Result.CopyValue(Inputs[i], Idxs[i][j].FromIdx, Idxs[i][j].ToIdx); end ); end; diff --git a/Src/Myc.Data.Records.pas b/Src/Myc.Data.Records.pas index 701c3c8..fa35036 100644 --- a/Src/Myc.Data.Records.pas +++ b/Src/Myc.Data.Records.pas @@ -4,8 +4,6 @@ interface uses System.SysUtils, - System.Generics.Collections, - System.Rtti, System.TypInfo; type @@ -23,8 +21,11 @@ type procedure FromType(const [ref] Buffer: TBytes; SrcType: PTypeInfo; const Src); procedure ToType(const [ref] Buffer: TBytes; DstType: PTypeInfo; var Dst); - procedure Assign(const [ref] Buffer: TBytes; const Src); overload; - procedure Finalize(const [ref] Buffer: TBytes); + procedure InitField(const [ref] Buffer: TBytes); + procedure AssignField(const [ref] Dest: TBytes; const [ref] Source: TBytes); + procedure FinalizeField(const [ref] Buffer: TBytes); + + procedure CopyField(Src, Dst: Pointer); property Offset: Integer read FOffset; @@ -51,6 +52,7 @@ type FFields: TArray; constructor Create(const AFields: TArray); + public class function FromRecord: TLayout; static; class function Construct(const Def: TArray): TLayout; static; @@ -60,15 +62,17 @@ type end; const - Align = 8; + Align = sizeof(Pointer); private FLayout: TLayout; FBuffer: TBytes; public - constructor Create(const ALayout: TLayout; const ABuffer: TBytes = nil); + constructor Create(const ALayout: TLayout); + class operator Finalize(var Dest: TDataRecord); + class operator Assign(var Dest: TDataRecord; const [ref] Src: TDataRecord); class function FromRecord: TDataRecord; overload; static; class function FromRecord(const Src: T): TDataRecord; overload; static; @@ -79,52 +83,54 @@ type function GetValue(const Name: String): T; overload; procedure GetValue(Idx: Integer; out Value); overload; - procedure CopyField(Idx: Integer; Dst: TDataRecord; DstIdx: Integer); + procedure CopyValue(const SrcRec: TDataRecord; SrcIdx, DstIdx: Integer); property Layout: TLayout read FLayout; end; +implementation + +uses + System.Rtti; + const DataSize: array[TDataRecord.TFieldType] of Integer = (sizeof(Double), sizeof(Int64), sizeof(String), sizeof(TDateTime), sizeof(TDataRecord)); -implementation +{ TDataRecord } -uses - System.Classes; - -constructor TDataRecord.Create(const ALayout: TLayout; const ABuffer: TBytes = nil); +constructor TDataRecord.Create(const ALayout: TLayout); begin FLayout := ALayout; - FBuffer := ABuffer; var bufSize := 0; if Length(FLayout.Fields) > 0 then with FLayout.Fields[High(FLayout.Fields)] do bufSize := Offset + AlignedSize; - if FBuffer = nil then - SetLength(FBuffer, bufSize) - else - Assert(Length(FBuffer) >= bufSize); + SetLength(FBuffer, bufSize); + for var i := 0 to High(FLayout.Fields) do + FLayout.Fields[i].InitField(FBuffer); end; -procedure TDataRecord.CopyField(Idx: Integer; Dst: TDataRecord; DstIdx: Integer); +procedure TDataRecord.CopyValue(const SrcRec: TDataRecord; SrcIdx, DstIdx: Integer); begin - Assert(Dst.Layout.Fields[DstIdx].FieldType = FLayout.Fields[Idx].FieldType); - Dst.Layout.Fields[DstIdx].Assign(Dst.FBuffer, FBuffer[FLayout.Fields[DstIdx].Offset]) -end; + Assert(SrcRec.Layout.Fields[SrcIdx].FieldType = FLayout.Fields[SrcIdx].FieldType); -{ TDataRecord } + var Src := @SrcRec.FBuffer[SrcRec.Layout.Fields[SrcIdx].Offset]; + var Dst := @FBuffer[FLayout.Fields[DstIdx].Offset]; + + FLayout.Fields[SrcIdx].CopyField(Src, Dst); +end; class function TDataRecord.FromRecord: TDataRecord; begin - Result := TDataRecord.Create(TLayout.FromRecord); + Result.Create(TLayout.FromRecord); end; class function TDataRecord.FromRecord(const Src: T): TDataRecord; begin - Result := FromRecord; + Result.Create(TLayout.FromRecord); var ctx := TRttiContext.Create; var rttiType := ctx.GetType(TypeInfo(T)); @@ -187,10 +193,22 @@ begin FLayout.Fields[idx].FromType(FBuffer, TypeInfo(T), Value); end; +class operator TDataRecord.Assign(var Dest: TDataRecord; const [ref] Src: TDataRecord); +begin + if Dest.FLayout.FFields <> Src.FLayout.FFields then + begin + Finalize(Dest); + Dest.Create(Src.Layout); + end; + + for var i := 0 to High(Dest.Layout.Fields) do + Dest.Layout.Fields[i].AssignField(Dest.FBuffer, Src.FBuffer); +end; + class operator TDataRecord.Finalize(var Dest: TDataRecord); begin - for var i := 0 to High(Dest.FLayout.Fields) do - Dest.FLayout.Fields[i].Finalize(Dest.FBuffer); + for var i := High(Dest.FLayout.Fields) downto 0 do + Dest.Layout.Fields[i].FinalizeField(Dest.FBuffer); end; constructor TDataRecord.TField.Create(const AName: string; AFieldType: TFieldType; AOffset: Integer); @@ -200,7 +218,7 @@ begin FOffset := AOffset; end; -procedure TDataRecord.TField.Finalize(const [ref] Buffer: TBytes); +procedure TDataRecord.TField.InitField(const [ref] Buffer: TBytes); begin if not (FFieldType in [dfString, dfRecord]) then exit; @@ -209,8 +227,22 @@ begin var P := @Buffer[FOffset]; case FFieldType of - dfString: PString(P)^ := ''; - dfRecord: TDataRecord(P^) := Default(TDataRecord); + dfString: Initialize(PString(P)^); + dfRecord: Initialize(TDataRecord(P^)); + end; +end; + +procedure TDataRecord.TField.FinalizeField(const [ref] Buffer: TBytes); +begin + if not (FFieldType in [dfString, dfRecord]) then + exit; + + Assert(FOffset + Size <= Length(Buffer)); + + var P := @Buffer[FOffset]; + case FFieldType of + dfString: Finalize(PString(P)^); + dfRecord: Finalize(TDataRecord(P^)); end; end; @@ -218,8 +250,6 @@ procedure TDataRecord.TField.FromType(const [ref] Buffer: TBytes; SrcType: PType begin Assert(FOffset + Size <= Length(Buffer)); - Finalize(Buffer); - var Dst := @Buffer[FOffset]; case FFieldType of dfFloat: @@ -317,19 +347,24 @@ begin end; end; -procedure TDataRecord.TField.Assign(const [ref] Buffer: TBytes; const Src); +procedure TDataRecord.TField.AssignField(const [ref] Dest: TBytes; const [ref] Source: TBytes); begin - Assert(FOffset + Size <= Length(Buffer)); + Assert(FOffset + Size <= Length(Dest)); - Finalize(Buffer); + var Dst := @Dest[FOffset]; + var Src := @Source[FOffset]; - var Dst := @Buffer[FOffset]; + CopyField(Src, Dst); +end; + +procedure TDataRecord.TField.CopyField(Src, Dst: Pointer); +begin case FFieldType of - dfFloat: PDouble(Dst)^ := PDouble(@Src)^; - dfInteger: PInt64(Dst)^ := PInt64(@Src)^; - dfString: PString(Dst)^ := PString(@Src)^; - dfTimestamp: PDateTime(Dst)^ := PDateTime(@Src)^; - dfRecord: TDataRecord(Dst^) := TDataRecord(Src); + dfFloat: PDouble(Dst)^ := PDouble(Src)^; + dfInteger: PInt64(Dst)^ := PInt64(Src)^; + dfString: PString(Dst)^ := PString(Src)^; + dfTimestamp: PDateTime(Dst)^ := PDateTime(Src)^; + dfRecord: TDataRecord(Dst^) := TDataRecord(Src^); else Assert(false); end; diff --git a/Src/Myc.Trade.Indicators.pas b/Src/Myc.Trade.Indicators.pas index 04feb51..e4e6140 100644 --- a/Src/Myc.Trade.Indicators.pas +++ b/Src/Myc.Trade.Indicators.pas @@ -8,8 +8,9 @@ uses System.Generics.Collections, System.Rtti, Myc.Data.Records, - Myc.Trade.Types, - Myc.Data.Series; + Myc.Data.Pipeline, + Myc.Data.Series, + Myc.Trade.Types; type // Result for the Moving Average Convergence Divergence (MACD) indicator. @@ -46,48 +47,48 @@ type class function CalculateWMA(const Series: TSeries; const Period: Integer): Double; static; public // Simple Moving Average - class function CreateSMA(Period: Integer): TConstFunc; static; + class function CreateSMA(Period: Integer): TConvertFunc; static; // Exponential Moving Average - class function CreateEMA(Period: Integer): TConstFunc; static; + class function CreateEMA(Period: Integer): TConvertFunc; static; // Hull Moving Average - class function CreateHMA(Period: Integer): TConstFunc; static; + class function CreateHMA(Period: Integer): TConvertFunc; static; // Relative Strength Index - class function CreateRSI(Period: Integer): TConstFunc; static; + class function CreateRSI(Period: Integer): TConvertFunc; static; // Moving Average Convergence Divergence - class function CreateMACD(FastPeriod, SlowPeriod, SignalPeriod: Integer): TConstFunc; overload; static; + class function CreateMACD(FastPeriod, SlowPeriod, SignalPeriod: Integer): TConvertFunc; overload; static; class function CreateMACD( const EmaFast, EmaSlow, - EmaSignal: TConstFunc - ): TConstFunc; overload; static; + EmaSignal: TConvertFunc + ): TConvertFunc; overload; static; // Stochastic Oscillator - class function CreateStochastic(KPeriod, DPeriod: Integer): TConstFunc; overload; static; + class function CreateStochastic(KPeriod, DPeriod: Integer): TConvertFunc; overload; static; class function CreateStochastic( KPeriod: Integer; - const SmaD: TConstFunc - ): TConstFunc; overload; static; + const SmaD: TConvertFunc + ): TConvertFunc; overload; static; // Bollinger Bands - class function CreateBollingerBands(Period: Integer; Multiplier: Double): TConstFunc; static; + class function CreateBollingerBands(Period: Integer; Multiplier: Double): TConvertFunc; static; // Average True Range - class function CreateATR(Period: Integer): TConstFunc; overload; static; - class function CreateATR(const MovAvgTR: TConstFunc): TConstFunc; overload; static; + class function CreateATR(Period: Integer): TConvertFunc; overload; static; + class function CreateATR(const MovAvgTR: TConvertFunc): TConvertFunc; overload; static; // Keltner Channels class function CreateKeltnerChannels( Period: Integer; Multiplier: Double - ): TConstFunc; overload; static; + ): TConvertFunc; overload; static; class function CreateKeltnerChannels( - const MovAvgMiddle: TConstFunc; - const AtrFunc: TConstFunc; + const MovAvgMiddle: TConvertFunc; + const AtrFunc: TConvertFunc; Multiplier: Double - ): TConstFunc; overload; static; + ): TConvertFunc; overload; static; - class function CreateMean: TConstFunc, Double>; static; + class function CreateMean: TConvertFunc, Double>; static; end; TIndicatorFactory = class type - TFunc = TConstFunc; + TFunc = TConvertFunc; private FParams: TDataRecord.TLayout; FInput: TDataRecord.TLayout; @@ -125,9 +126,9 @@ type TMACD = class type TParam = record - Fast: TConstFunc; - Slow: TConstFunc; - Signal: TConstFunc; + Fast: TConvertFunc; + Slow: TConvertFunc; + Signal: TConvertFunc; end; TInput = record @@ -141,7 +142,7 @@ type end; public - class function CreateMACD(const Param: TParam): TConstFunc; static; + class function CreateMACD(const Param: TParam): TConvertFunc; static; end; var @@ -208,7 +209,7 @@ begin Result := numerator / denominator; end; -class function TIndicators.CreateBollingerBands(Period: Integer; Multiplier: Double): TConstFunc; +class function TIndicators.CreateBollingerBands(Period: Integer; Multiplier: Double): TConvertFunc; begin var sourceData: TSeries; Result := @@ -231,7 +232,7 @@ begin end; end; -class function TIndicators.CreateEMA(Period: Integer): TConstFunc; +class function TIndicators.CreateEMA(Period: Integer): TConvertFunc; begin var lastEma: Double := Double.NaN; var sourceData: TSeries; @@ -262,7 +263,7 @@ begin end; end; -class function TIndicators.CreateHMA(Period: Integer): TConstFunc; +class function TIndicators.CreateHMA(Period: Integer): TConvertFunc; begin var periodHalf := Period div 2; var periodSqrt := Round(Sqrt(Period)); @@ -305,13 +306,13 @@ begin end; // Standard MACD using EMAs. -class function TIndicators.CreateMACD(FastPeriod, SlowPeriod, SignalPeriod: Integer): TConstFunc; +class function TIndicators.CreateMACD(FastPeriod, SlowPeriod, SignalPeriod: Integer): TConvertFunc; begin Result := CreateMACD(CreateEMA(FastPeriod), CreateEMA(SlowPeriod), CreateEMA(SignalPeriod)); end; // Creates a MACD indicator from three provided moving average functions. -class function TIndicators.CreateMACD(const EmaFast, EmaSlow, EmaSignal: TConstFunc): TConstFunc; +class function TIndicators.CreateMACD(const EmaFast, EmaSlow, EmaSignal: TConvertFunc): TConvertFunc; begin Result := function(const Value: Double): TMacdResult @@ -339,7 +340,7 @@ begin end; end; -class function TIndicators.CreateRSI(Period: Integer): TConstFunc; +class function TIndicators.CreateRSI(Period: Integer): TConvertFunc; begin var avgGain: Double := Double.NaN; var avgLoss: Double := Double.NaN; @@ -398,7 +399,7 @@ begin end; end; -class function TIndicators.CreateSMA(Period: Integer): TConstFunc; +class function TIndicators.CreateSMA(Period: Integer): TConvertFunc; begin var sourceData: TSeries; Result := @@ -413,7 +414,7 @@ begin end; // Standard Stochastic Oscillator using an SMA for the %D line. -class function TIndicators.CreateStochastic(KPeriod, DPeriod: Integer): TConstFunc; +class function TIndicators.CreateStochastic(KPeriod, DPeriod: Integer): TConvertFunc; begin Result := CreateStochastic(KPeriod, CreateSMA(DPeriod)); end; @@ -421,8 +422,8 @@ end; // Creates a Stochastic Oscillator using an injectable moving average for the %D line. class function TIndicators.CreateStochastic( KPeriod: Integer; - const SmaD: TConstFunc -): TConstFunc; + const SmaD: TConvertFunc +): TConvertFunc; begin var sourceData: TSeries; @@ -461,13 +462,13 @@ begin end; // Standard ATR using an EMA for smoothing. -class function TIndicators.CreateATR(Period: Integer): TConstFunc; +class function TIndicators.CreateATR(Period: Integer): TConvertFunc; begin Result := CreateATR(CreateEMA(Period)); end; // Calculates the Average True Range (ATR) using an injectable moving average. -class function TIndicators.CreateATR(const MovAvgTR: TConstFunc): TConstFunc; +class function TIndicators.CreateATR(const MovAvgTR: TConvertFunc): TConvertFunc; begin var sourceData: TSeries; @@ -495,17 +496,17 @@ begin end; // Standard Keltner Channels using an EMA for the middle line and an EMA-based ATR. -class function TIndicators.CreateKeltnerChannels(Period: Integer; Multiplier: Double): TConstFunc; +class function TIndicators.CreateKeltnerChannels(Period: Integer; Multiplier: Double): TConvertFunc; begin Result := CreateKeltnerChannels(CreateEMA(Period), CreateATR(Period), Multiplier); end; // Calculates Keltner Channels using an injectable ATR and middle band moving average. class function TIndicators.CreateKeltnerChannels( - const MovAvgMiddle: TConstFunc; - const AtrFunc: TConstFunc; + const MovAvgMiddle: TConvertFunc; + const AtrFunc: TConvertFunc; Multiplier: Double -): TConstFunc; +): TConvertFunc; begin Result := function(const Value: TOhlcItem): TKeltnerChannelsResult @@ -533,7 +534,7 @@ begin end; end; -class function TIndicators.CreateMean: TConstFunc, Double>; +class function TIndicators.CreateMean: TConvertFunc, Double>; begin Result := function(const Value: TArray): Double @@ -548,7 +549,7 @@ begin end; // Creates a MACD indicator from three provided moving average functions. -class function TMACD.CreateMACD(const Param: TParam): TConstFunc; +class function TMACD.CreateMACD(const Param: TParam): TConvertFunc; begin Result := function(const Input: TInput): TResult diff --git a/Src/Myc.Trade.Pipeline.Impl.pas b/Src/Myc.Trade.Pipeline.Impl.pas new file mode 100644 index 0000000..b14a3bf --- /dev/null +++ b/Src/Myc.Trade.Pipeline.Impl.pas @@ -0,0 +1,218 @@ +unit Myc.Trade.Pipeline.Impl; + +interface + +uses + Myc.Signals, + Myc.Data.Pipeline, + Myc.Data.Pipeline.Impl, + Myc.Trade.Types; + +type + TTickAggregation = class(TMycConverter, TDataPoint>) + private + FTimeframe: TTimeframe; + FCurrentBar: TDataPoint; + 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); + property CurrentBar: TDataPoint read GetCurrentBar; + property Timeframe: TTimeframe read GetTimeframe; + end; + + TOhlcAggregation = class(TMycConverter, TDataPoint>) + private + FTimeframe: TTimeframe; + FCurrentBar: TDataPoint; + 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); + property CurrentBar: TDataPoint read GetCurrentBar; + property Timeframe: TTimeframe read GetTimeframe; + end; + +implementation + +uses + System.Math, + System.DateUtils; + +{ TTickAggregation } + +constructor TTickAggregation.Create(const ATimeframe: TTimeframe); +begin + inherited Create; + FTimeframe := ATimeframe; +end; + +function TTickAggregation.GetBarStartTime(const TimeStamp: TDateTime; const Timeframe: TTimeframe): TDateTime; +var + baseTime: TDateTime; +begin + // 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)); + M5: Result := RecodeMinute(RecodeSecond(baseTime, 0), MinuteOf(TimeStamp) - (MinuteOf(TimeStamp) mod 5)); + 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: Result := Floor(TimeStamp) - (Floor(TimeStamp) mod 2); + D3: Result := Floor(TimeStamp) - (Floor(TimeStamp) mod 3); + W: Result := TimeStamp.StartOfTheWeek; + MN: Result := TimeStamp.StartOfTheMonth; + MN3: Result := RecodeMonth(TimeStamp.StartOfTheMonth, (MonthOf(TimeStamp) - 1) div 3 * 3 + 1); + MN6: Result := RecodeMonth(TimeStamp.StartOfTheMonth, (MonthOf(TimeStamp) - 1) div 6 * 6 + 1); + Y: Result := TimeStamp.StartOfTheYear; + else + Result := 0; + end; +end; + +function TTickAggregation.GetCurrentBar: TDataPoint; +begin + Result := FCurrentBar; +end; + +function TTickAggregation.GetTimeframe: TTimeframe; +begin + Result := FTimeframe; +end; + +function TTickAggregation.Consume(const Value: TDataPoint): TState; +var + barStartTime: TDateTime; + lastBarTime: TDateTime; +begin + barStartTime := GetBarStartTime(Value.Time, FTimeframe); + lastBarTime := FCurrentBar.Time; + + if (barStartTime > lastBarTime) then + begin + if (lastBarTime > 0) then + begin + Result := Broadcast(FCurrentBar); + end; + + FCurrentBar.Data.Open := Value.Data; + FCurrentBar.Data.High := Value.Data; + FCurrentBar.Data.Low := Value.Data; + FCurrentBar.Data.Close := Value.Data; + FCurrentBar.Data.Volume := 1; + FCurrentBar.Time := barStartTime; + end + else + begin + 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; + FCurrentBar.Data.Volume := FCurrentBar.Data.Volume + 1; + end; +end; + +{ TOhlcAggregation } + +constructor TOhlcAggregation.Create(const ATimeframe: TTimeframe); +begin + inherited Create; + FTimeframe := ATimeframe; +end; + +function TOhlcAggregation.GetBarStartTime(const TimeStamp: TDateTime; const Timeframe: TTimeframe): TDateTime; +begin + // 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)); + M5: Result := RecodeMinute(RecodeSecond(baseTime, 0), MinuteOf(TimeStamp) - (MinuteOf(TimeStamp) mod 5)); + 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: Result := Floor(TimeStamp) - (Floor(TimeStamp) mod 2); + D3: Result := Floor(TimeStamp) - (Floor(TimeStamp) mod 3); + W: Result := TimeStamp.StartOfTheWeek; + MN: Result := TimeStamp.StartOfTheMonth; + MN3: Result := RecodeMonth(TimeStamp.StartOfTheMonth, (MonthOf(TimeStamp) - 1) div 3 * 3 + 1); + MN6: Result := RecodeMonth(TimeStamp.StartOfTheMonth, (MonthOf(TimeStamp) - 1) div 6 * 6 + 1); + Y: Result := TimeStamp.StartOfTheYear; + else + Result := 0; + end; +end; + +function TOhlcAggregation.GetCurrentBar: TDataPoint; +begin + Result := FCurrentBar; +end; + +function TOhlcAggregation.GetTimeframe: TTimeframe; +begin + Result := FTimeframe; +end; + +function TOhlcAggregation.Consume(const Value: TDataPoint): TState; +var + barStartTime: TDateTime; + lastBarTime: TDateTime; +begin + barStartTime := GetBarStartTime(Value.Time, FTimeframe); + lastBarTime := FCurrentBar.Time; + + if (barStartTime > lastBarTime) then + begin + if (lastBarTime > 0) then + begin + Result := Broadcast(FCurrentBar); + end; + + FCurrentBar.Data := Value.Data; + FCurrentBar.Time := barStartTime; + end + else + begin + 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; + FCurrentBar.Data.Volume := FCurrentBar.Data.Volume + Value.Data.Volume; + end; +end; + +end. diff --git a/Src/Myc.Trade.Pipeline.pas b/Src/Myc.Trade.Pipeline.pas new file mode 100644 index 0000000..8319dd6 --- /dev/null +++ b/Src/Myc.Trade.Pipeline.pas @@ -0,0 +1,30 @@ +unit Myc.Trade.Pipeline; + +interface + +uses + Myc.Data.Pipeline, + Myc.Trade.Types; + +type + TTradeConverter = record + class function CreateTickAggregation(Timeframe: TTimeframe): TConverter, TDataPoint>; static; + class function CreateOhlcAggregation(Timeframe: TTimeframe): TConverter, TDataPoint>; static; + end; + +implementation + +uses + Myc.Trade.Pipeline.Impl; + +class function TTradeConverter.CreateTickAggregation(Timeframe: TTimeframe): TConverter, TDataPoint>; +begin + Result := TTickAggregation.Create(Timeframe); +end; + +class function TTradeConverter.CreateOhlcAggregation(Timeframe: TTimeframe): TConverter, TDataPoint>; +begin + Result := TOhlcAggregation.Create(Timeframe); +end; + +end. diff --git a/Src/Myc.Trade.Types.pas b/Src/Myc.Trade.Types.pas index eb0282e..7a51fb6 100644 --- a/Src/Myc.Trade.Types.pas +++ b/Src/Myc.Trade.Types.pas @@ -2,6 +2,9 @@ unit Myc.Trade.Types; interface +uses + Myc.Core.Future; + type TTimeframe = (S, S5, S15, S30, M, M2, M3, M5, M10, M15, M30, H, H2, H3, H4, H8, H12, D, D2, D3, W, MN, MN3, MN6, Y); @@ -28,11 +31,6 @@ type constructor Create(ATime: TDateTime; const AData: T); end; - TConstFunc = reference to function(const Value: S): T; - TConstFunc = reference to function(const Value1: S; const Value2: T): U; - TConstProc = reference to procedure(const Value: T); - TConstFuncPredicate = reference to function(const Value: S; out Res: T): Boolean; - implementation { TAskBidItem } diff --git a/Test/TestDataArray.pas b/Test/TestDataArray.pas index ab85bed..baf7095 100644 --- a/Test/TestDataArray.pas +++ b/Test/TestDataArray.pas @@ -5,7 +5,7 @@ interface uses System.SysUtils, DUnitX.TestFramework, - Myc.Trade.DataArray; + Myc.Data.Series; type [TestFixture] diff --git a/Test/TestDataRecord.pas b/Test/TestDataRecord.pas index 3aac302..18a017f 100644 --- a/Test/TestDataRecord.pas +++ b/Test/TestDataRecord.pas @@ -7,7 +7,7 @@ uses System.SysUtils, System.DateUtils, System.Rtti, - Myc.DataRecord; + Myc.Data.Records; type // A record that contains all supported data types for testing @@ -52,6 +52,10 @@ type // [IgnoreMemoryLeaks] procedure TestNestedDataRecord; + [Test] + // [IgnoreMemoryLeaks] + procedure TestInternalRecordAssignment; + [Test] // [IgnoreMemoryLeaks] procedure TestNestedSetValueAndGetValue; @@ -174,7 +178,8 @@ var begin // Setup: Create the nested record and convert it to a TDataRecord nestedSrc.NestedValue := 999; - nestedSrc.NestedString := 'Nested Record Test'; + const str = 'Nested Record Test'; + nestedSrc.NestedString := str; nestedDataRecord := TDataRecord.FromRecord(nestedSrc); // Setup: Create the outer record containing the nested TDataRecord @@ -192,12 +197,41 @@ begin // Assert values within the nested record Assert.AreEqual(999, retrievedNestedRecord.GetValue('NestedValue'), 'Nested integer value mismatch'); - Assert.AreEqual('Nested Record Test', retrievedNestedRecord.GetValue('NestedString'), 'Nested string value mismatch'); + Assert.AreEqual(str, retrievedNestedRecord.GetValue('NestedString'), 'Nested string value mismatch'); // Assert value in the outer record Assert.AreEqual(111, mainDataRecord.GetValue('OuterValue'), 'Outer integer value mismatch'); end; +procedure TTestDataRecord.TestInternalRecordAssignment; +var + nestedSrc: TNestedTestRecord; + nestedDataRecord: TDataRecord; + outerSrc: TOuterTestRecord; + mainDataRecord: TDataRecord; + retrievedNestedRecord: TDataRecord; +begin + // Setup: Create the nested record and convert it to a TDataRecord + const str = 'Nested Record Test'; + nestedSrc.NestedString := str; + nestedDataRecord := TDataRecord.FromRecord(nestedSrc); + + // Setup: Create the outer record containing the nested TDataRecord + outerSrc.MyNestedRecord := nestedDataRecord; + + // Place a record on the heap. FAILS: This erases the nested record data!! + TDataRecord.FromRecord(outerSrc); + + // Test: Create the main TDataRecord from the outer record instance + mainDataRecord := TDataRecord.FromRecord(outerSrc); + + // Assert: Retrieve the nested record and check its contents + retrievedNestedRecord := mainDataRecord.GetValue('MyNestedRecord'); + + // Assert values within the nested record + Assert.AreEqual(str, retrievedNestedRecord.GetValue('NestedString'), 'Nested string value mismatch'); +end; + procedure TTestDataRecord.TestNestedSetValueAndGetValue; var outerRecord: TDataRecord;