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.