diff --git a/Src/Myc.Trade.DataPoint.pas b/Src/Myc.Trade.DataPoint.pas index 7a08879..a5be20b 100644 --- a/Src/Myc.Trade.DataPoint.pas +++ b/Src/Myc.Trade.DataPoint.pas @@ -13,19 +13,21 @@ type end; TDataPoint = record - Idx: Int64; Time: TDateTime; Data: T; - constructor Create(AIdx: Int64; ATime: TDateTime; const AData: T); + constructor Create(ATime: TDateTime; const AData: T); + end; + + IMycStream = interface + procedure Add(const Data: TDataPoint); end; implementation { TDataPoint } -constructor TDataPoint.Create(AIdx: Int64; ATime: TDateTime; const AData: T); +constructor TDataPoint.Create(ATime: TDateTime; const AData: T); begin - Idx := AIdx; Time := ATime; Data := AData; end; diff --git a/Src/Myc.Trade.DataStream.pas b/Src/Myc.Trade.DataStream.pas index 3959b46..74cbe0d 100644 --- a/Src/Myc.Trade.DataStream.pas +++ b/Src/Myc.Trade.DataStream.pas @@ -72,8 +72,8 @@ type IAuraDataServer = interface(IDataServer) ['{481E27DE-DC58-49AF-B8A8-B6316980C423}'] - function LoadDataFile(const FileName: string): TFuture>>; - function LoadDataSeries(const FileName: string): TFuture>>; + function LoadDataFile(const FileName: string): TFuture>>; + function LoadDataSeries(const FileName: string): TFuture>>; function EnumerateAssetFiles: TArray; end; @@ -88,7 +88,7 @@ type Name: String; Age: TDateTime; LastUsed: TDateTime; - Data: TFuture>>; + Data: TFuture>>; end; private // The cache is per-instance. @@ -102,8 +102,8 @@ type FLoadGate: TLatch; protected - class function ReadCompressedData(const InputStream: TStream): TArray>; static; - class function ReadUncompressedData(const InputStream: TStream): TArray>; static; + class function ReadCompressedData(const InputStream: TStream): TArray>; static; + class function ReadUncompressedData(const InputStream: TStream): TArray>; static; public constructor Create(const APath: String); @@ -124,8 +124,8 @@ type procedure ClearCache; function EnumerateAssetFiles: TArray; - function LoadDataFile(const FileName: string): TFuture>>; - function LoadDataSeries(const FileName: string): TFuture>>; + function LoadDataFile(const FileName: string): TFuture>>; + function LoadDataSeries(const FileName: string): TFuture>>; property Path: String read GetPath; end; @@ -136,11 +136,10 @@ type FDataServer: TAuraDataServer; FIsLiveData: TMutable.IWriteable; FCurrentFileName: string; - FCurrentData: TFuture>>; + FCurrentData: TFuture>>; FNextFileName: string; - FNextData: TFuture>>; + FNextData: TFuture>>; FCurrPosInFile: Int64; - FCurrentIdx: Int64; FLastTimeStamp: TDateTime; protected function GetChunk(var Data: array of TDataPoint): Integer; override; @@ -315,12 +314,12 @@ begin Result := FPath; end; -function TAuraDataServer.LoadDataFile(const FileName: string): TFuture>>; +function TAuraDataServer.LoadDataFile(const FileName: string): TFuture>>; var parsedPath, parsedSymbol: string; parsedYear, parsedMonth: Integer; begin - Result := TFuture>>.Null; + Result := TFuture>>.Null; if not TryParseFileName(FileName, parsedPath, parsedSymbol, parsedYear, parsedMonth) then exit; @@ -369,8 +368,8 @@ begin .Construct( TLatch.Enqueue(FLoadGate), // Use shared class var FLoadGate function: TBytes begin Result := TFile.ReadAllBytes(capFileName); end) - .Chain>>( - function(bytes: TBytes): TArray> + .Chain>>( + function(bytes: TBytes): TArray> begin var zipMemoryStream := TBytesStream.Create(bytes); try @@ -384,9 +383,9 @@ begin else if TFile.Exists(FileName) then begin Result := - TFuture>>.Construct( + TFuture>>.Construct( TLatch.Enqueue(FLoadGate), // Use shared class var FLoadGate - function: TArray> + function: TArray> begin if TFile.Exists(capFileName) then begin @@ -419,11 +418,11 @@ begin end; end; -function TAuraDataServer.LoadDataSeries(const FileName: string): TFuture>>; +function TAuraDataServer.LoadDataSeries(const FileName: string): TFuture>>; var loadedState: TState; - loadedFiles: TArray>>>; - liveData: TFuture>>; + loadedFiles: TArray>>>; + liveData: TFuture>>; tabFiles: TList; currentFile: string; liveFilePath: string; @@ -453,7 +452,7 @@ begin end; end; - var loadedFileList := TList>>>.Create; + var loadedFileList := TList>>>.Create; var loadStates := TList.Create; try for currentFile in tabFiles do @@ -463,7 +462,7 @@ begin loadStates.Add(data.Done); end; - liveData := TFuture>>.Null; + liveData := TFuture>>.Null; if liveFilePath <> '' then begin liveData := LoadDataFile(liveFilePath); @@ -482,11 +481,11 @@ begin end; Result := - TFuture>>.Construct( + TFuture>>.Construct( loadedState, - function: TArray> + function: TArray> begin - var tickList := TList>.Create; + var tickList := TList>.Create; try var overallLastTabTickTime: TDateTime := 0; var cnt := 0; @@ -496,10 +495,10 @@ begin for var P in loadedFiles do for var iterTickData in P.Value do - if overallLastTabTickTime < iterTickData.TimeStamp then + if overallLastTabTickTime < iterTickData.Time then begin tickList.Add(iterTickData); - overallLastTabTickTime := iterTickData.TimeStamp; + overallLastTabTickTime := iterTickData.Time; end; Result := tickList.ToArray; finally @@ -509,7 +508,7 @@ begin ); end; -class function TAuraDataServer.ReadCompressedData(const InputStream: TStream): TArray>; +class function TAuraDataServer.ReadCompressedData(const InputStream: TStream): TArray>; var decompressionStream: TStream; localHeader: TZipHeader; @@ -548,10 +547,11 @@ begin end; end; -class function TAuraDataServer.ReadUncompressedData(const InputStream: TStream): TArray>; +class function TAuraDataServer.ReadUncompressedData(const InputStream: TStream): TArray>; var fileSize: Int64; recordCount, bytesRead: Integer; + rec: TAuraFileDataRecord; begin SetLength(Result, 0); InputStream.Position := 0; @@ -563,9 +563,13 @@ begin if recordCount > 0 then begin SetLength(Result, recordCount); - bytesRead := InputStream.Read(Result[0], fileSize); - if bytesRead <> fileSize then - raise EReadError.CreateFmt('Read error. Expected %d bytes, read %d.', [fileSize, bytesRead]); + for var i:=0 to High(Result) do + begin + bytesRead := InputStream.Read(rec, SizeOf(TAuraFileDataRecord)); + if bytesRead <> SizeOf(TAuraFileDataRecord) then + raise EReadError.CreateFmt('Read error. Expected %d bytes, read %d.', [fileSize, bytesRead]); + Result[i].Create( rec.TimeStamp, rec.Data ); + end; end; end; @@ -585,7 +589,6 @@ begin inherited; FCurrentData := FDataServer.LoadDataFile(FCurrentFileName); FCurrPosInFile := 0; - FCurrentIdx := 0; FIsLiveData.SetValue(false); FLastTimeStamp := 0; FNextFileName := FDataServer.FindNextDataFile(FCurrentFileName); @@ -597,8 +600,8 @@ end; function TAuraFileStream.GetChunk(var Data: array of TDataPoint): Integer; var - item: TAuraFileDataRecord; - currData: TArray>; + item: TDataPoint; + currData: TArray>; begin Result := 0; if not FCurrentData.Done.IsSet then @@ -635,12 +638,11 @@ begin end; item := currData[FCurrPosInFile]; - if FLastTimeStamp < item.TimeStamp then + if FLastTimeStamp < item.Time then begin - FLastTimeStamp := item.TimeStamp; - Data[Result].Create(FCurrentIdx, FLastTimeStamp, item.Data); + FLastTimeStamp := item.Time; + Data[Result].Create(FLastTimeStamp, item.Data); Inc(Result); - Inc(FCurrentIdx); end; Inc(FCurrPosInFile); end; diff --git a/Test/MycTests.res b/Test/MycTests.res index 333684a..e30ea80 100644 Binary files a/Test/MycTests.res and b/Test/MycTests.res differ