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