unit Myc.Trade.DataProvider; interface uses Myc.Signals, Myc.Lazy, Myc.Trade.DataPoint, Myc.Trade.DataStream; type TDataStreamProvider = class(TInterfacedObject, TMutable>.IMutable) private FStream: IDataStream; FDataArray: IDataArray; FNewDataAvailable: TFlag; FChunk: TArray>; FLookback: Int64; function GetChanged: TSignal; function GetValue: TDataSeries; public constructor Create(ALookback, AMaxChunkSize: Int64; const AStream: IDataStream); destructor Destroy; override; end; TDataStreamProvider = record class function Create(Lookback, MaxChunkSize: Int64; const Stream: IDataStream): TMutable>; static; end; implementation { TDataStreamProvider } constructor TDataStreamProvider.Create(ALookback, AMaxChunkSize: Int64; const AStream: IDataStream); begin inherited Create; FLookback := ALookback; SetLength(FChunk, AMaxChunkSize); FStream := AStream; FNewDataAvailable := TFlag.CreateObserver(FStream.HasData); FDataArray := TDataSeries.CreateArray(ALookback); end; destructor TDataStreamProvider.Destroy; begin inherited; end; function TDataStreamProvider.GetChanged: TSignal; begin Result := FStream.HasData; end; function TDataStreamProvider.GetValue: TDataSeries; begin if FNewDataAvailable.Reset then begin var cnt := FStream.GetChunk(FChunk); if cnt > 0 then FDataArray.Add(FChunk, cnt); end; Result := FDataArray; end; { TDataStreamProvider } class function TDataStreamProvider.Create(Lookback, MaxChunkSize: Int64; const Stream: IDataStream): TMutable>; begin Result := TDataStreamProvider.Create(Lookback, MaxChunkSize, Stream); end; end.