diff --git a/AuraTrader/MainForm.pas b/AuraTrader/MainForm.pas index d766d67..bc1359e 100644 --- a/AuraTrader/MainForm.pas +++ b/AuraTrader/MainForm.pas @@ -102,11 +102,8 @@ type private FOnEvent: TNotifyEvent; { Private declarations } -{$ifdef TICKDATA} - FServer: IDataServer; -{$else} - FServer: IDataServer; -{$endif} + FFileCache: IDataFileCache>; + FServer: IDataServer; FSymbols: TFuture>; FTerminate: TEvent; FProcessDone: TState; @@ -228,10 +225,11 @@ begin FTerminate := TEvent.CreateEvent; // Create an instance of the TAuraTABFileServer. The server can be reused for multiple stream creations. [364] + FFileCache := TDataFileCache>.Create; {$ifdef TICKDATA} - FServer := TAskBidFileServer.Create('\\COFFEE\TickData\Pepperstone'); + FServer := TTickFileServer.Create(FFileCache, '\\COFFEE\TickData\Pepperstone'); {$else} - FServer := TM1FileServer.Create('\\COFFEE\TickData\Pepperstone'); + FServer := TM1FileServer.Create(FFileCache, '\\COFFEE\TickData\Pepperstone'); {$endif} SymbolsComboBox.Enabled := false; @@ -339,7 +337,7 @@ begin var panel := pnlChart.AddPanel; panel.AddDoubleSeries(equity, TAlphaColors.Blue, 3); - FProcessDone := FProcessDone + FServer.ProcessData(Symbol, terminated, ticker.Consumer); + FProcessDone := FProcessDone + (FServer as IM1FileServer).ProcessData(Symbol, terminated, ticker.Consumer); end; @@ -587,44 +585,28 @@ begin var terminated := TFlag.CreateObserver(FTerminate.Signal).State; {$ifdef TICKDATA} - var ticker := TConverter.CreateTicker>; + var ticker := TConverter.CreateTicker>; var lastPrice := - ticker.Chain>( - function(const Tick: TDataPoint): TDataPoint + ticker.Producer.Chain>( + function(const Tick: TDataPoint): TDataPoint begin Result.Time := Tick.Time; Result.Data := 0.5 * (Tick.Data.Ask + Tick.Data.Bid); end ); - var OhlcPoint := lastPrice.Chain>(TConverter.CreateTickAggregation(Timeframe)); - OhlcPoint.Producer.Link(Consumer); + var OhlcPoint := lastPrice.Chain>(TTradeConverter.CreateTickAggregation(Timeframe)); + OhlcPoint.CreateLink(Consumer); - var Producer := - TConverter>, TArray>>.CreateGeneric( - function(const Values: TArray>): TArray> - begin - SetLength(Result, Length(Values)); - for var i := 0 to High(Result) do - begin - Result[i].Time := Values[i].Time; - Result[i].Data.Ask := Values[i].Data.Ask; - Result[i].Data.Bid := Values[i].Data.Bid; - end; - end - ); - - Producer.Producer.Link(ticker); - - FProcessDone := FProcessDone + FServer.ProcessData(Symbol, terminated, Producer); + FProcessDone := FProcessDone + (FServer as ITickFileServer).ProcessData(Symbol, terminated, ticker.Consumer); {$else} var ticker := TConverter.CreateTicker>; ticker.Producer.Chain(Consumer); - FProcessDone := FProcessDone + FServer.ProcessData(Symbol, terminated, ticker.Consumer); + FProcessDone := FProcessDone + (FServer as IM1FileServer).ProcessData(Symbol, terminated, ticker.Consumer); {$endif} end; diff --git a/Src/Myc.Core.FileCache.pas b/Src/Myc.Core.FileCache.pas index 53782a6..0d9cfe1 100644 --- a/Src/Myc.Core.FileCache.pas +++ b/Src/Myc.Core.FileCache.pas @@ -77,11 +77,12 @@ end; procedure TDataFileCache.Clear; begin + // Wait for any pending load operations to finish before clearing. + for var cachedFile in FCachedFiles.Values do + cachedFile.Data.WaitFor; + FLock.Enter; try - // Wait for any pending load operations to finish before clearing. - for var cachedFile in FCachedFiles.Values do - cachedFile.Data.WaitFor; FCachedFiles.Clear; finally FLock.Leave; diff --git a/Src/Myc.Trade.DataStream.pas b/Src/Myc.Trade.DataStream.pas index 6eff592..3a98484 100644 --- a/Src/Myc.Trade.DataStream.pas +++ b/Src/Myc.Trade.DataStream.pas @@ -34,29 +34,30 @@ type Bid: Double; end; - // Interface for an instantiable data server. - IDataServer = interface + IDataServer = interface procedure ClearCache; - function ProcessData(const Symbol: String; const Terminated: TState; const Processor: IConsumer>>): TState; function EnumerateSymbols: TFuture>; end; + // Interface for an instantiable data server. + IDataServer = interface(IDataServer) + function ProcessData(const Symbol: String; const Terminated: TState; const Processor: IConsumer>>): TState; + end; + // Metadata for a monthly data file (SYMBOL_YYYY_MM.EXT) TDataFile = record private - FExtension: String; FPath: String; - FSymbol: String; + FSymbol: String; // Includes extension, e.g. "EURUSD.M1" FYear: Integer; FMonth: Integer; FBasename: String; function GetIsValid: Boolean; public - constructor Create(const APath, ASymbol, AExtension: String; AYear, AMonth: Integer; const ABasename: String); + constructor Create(const APath, ASymbol: String; AYear, AMonth: Integer; const ABasename: String); function GetBaseFileName: string; function GetFullFileName: string; property IsValid: Boolean read GetIsValid; - property Extension: String read FExtension; property Path: String read FPath; property Symbol: String read FSymbol; property Year: Integer read FYear; @@ -64,22 +65,19 @@ type end; // Generic server base handling the pipeline and preloading. - TDataServer = class(TInterfacedObject, IDataServer) + TDataServer = class(TInterfacedObject, IDataServer, IDataServer) private - FCachedFiles: TDataFileCache>>>; + // Cache stores binary blobs of the MAPPED data (TDataPoint) + FCachedFiles: IDataFileCache>; FPath: String; FSymbols: TFuture>>; function GetPath: String; - function ProcessChunks( - const DataChunks: TArray>>; - Terminated: TState; - Processor: IConsumer>> - ): TState; + function ProcessChunks(const MappedBlobs: TArray; Terminated: TState; Processor: IConsumer>>): TState; function ProcessFile( FileInfo: TDataFile; - const DataFile: TFuture>>>; + const DataFile: TFuture>; const Terminated: TState; Processor: IConsumer>> ): TState; @@ -87,17 +85,22 @@ type strict private class var FLoadGate: TLatch; + class constructor CreateClass; protected // Maps the raw file record to the internal TDataPoint function MapRecord(const RecordBuffer: R): TDataPoint; virtual; abstract; - function DoLoad(const FileName: string; const Gate: TLatch): TFuture>>>; + + // Returns Future of Processed Blobs (Serialized TDataPoint) + function DoLoad(const FileName: string; const Gate: TLatch): TFuture>; + function GetNextFile(const Curr: TDataFile): TDataFile; virtual; - function ReadCompressedData(const InputStream: TStream): TArray>>; - function ReadUncompressedData(const InputStream: TStream): TArray>>; + function ReadCompressedData(const InputStream: TStream): TArray; + function ReadUncompressedData(const InputStream: TStream): TArray; public - constructor Create(const APath: String); + // Cache is injected. It stores TArray. + constructor Create(const ACache: IDataFileCache>; const APath: String); destructor Destroy; override; procedure AfterConstruction; override; procedure ClearCache; @@ -105,18 +108,27 @@ type function FindFirstFile(const Symbol: string): TFuture; function ParseFileName(const FileName: string): TDataFile; function EnumerateSymbols: TFuture>; - function LoadDataFile(const DataFile: TDataFile): TFuture>>>; + + function LoadDataFile(const DataFile: TDataFile): TFuture>; function ProcessData(const Symbol: String; const Terminated: TState; const Processor: IConsumer>>): TState; property Path: String read GetPath; end; - TTickFileServer = class(TDataServer) + ITickFileServer = interface(IDataServer) + ['{A6DE8B82-B0FB-48DB-A39A-CE91D861C5D2}'] + end; + + TTickFileServer = class(TDataServer, ITickFileServer) protected function MapRecord(const RecordBuffer: TTickRecord): TDataPoint; override; end; - TM1FileServer = class(TDataServer) + IM1FileServer = interface(IDataServer) + ['{7BB980D9-1969-49A4-991F-908F2DDAC1A0}'] + end; + + TM1FileServer = class(TDataServer, IM1FileServer) protected function MapRecord(const RecordBuffer: TM1Record): TDataPoint; override; end; @@ -135,11 +147,10 @@ const { TDataFile } -constructor TDataFile.Create(const APath, ASymbol, AExtension: String; AYear, AMonth: Integer; const ABasename: String); +constructor TDataFile.Create(const APath, ASymbol: String; AYear, AMonth: Integer; const ABasename: String); begin FPath := APath; FSymbol := ASymbol; - FExtension := AExtension; FYear := AYear; FMonth := AMonth; FBasename := ABasename; @@ -151,8 +162,11 @@ begin end; function TDataFile.GetFullFileName: string; +var + ext: string; begin - Result := TPath.Combine(FPath, FBasename + FExtension); + ext := TPath.GetExtension(FSymbol); + Result := TPath.Combine(FPath, FBasename + ext); end; function TDataFile.GetIsValid: Boolean; @@ -162,17 +176,20 @@ end; { TDataServer } -constructor TDataServer.Create(const APath: String); +class constructor TDataServer.CreateClass; +begin + FLoadGate := TLatch.CreateLatch(0); +end; + +constructor TDataServer.Create(const ACache: IDataFileCache>; const APath: String); begin inherited Create; FPath := APath; - FCachedFiles := TDataFileCache>>>.Create; + FCachedFiles := ACache; end; destructor TDataServer.Destroy; begin - ClearCache; - FCachedFiles.Free; inherited Destroy; end; @@ -184,7 +201,8 @@ end; procedure TDataServer.ClearCache; begin - FCachedFiles.Clear; + if Assigned(FCachedFiles) then + FCachedFiles.Clear; end; procedure TDataServer.UpdateSymbols; @@ -227,7 +245,7 @@ begin TComparer.Construct( function(const Left, Right: TDataFile): Integer begin - if Left.Year <> Right.Year then + if (Left.Year <> Right.Year) then Result := Left.Year - Right.Year else Result := Left.Month - Right.Month; @@ -264,7 +282,7 @@ begin var symbolFiles: TArray; begin - if Symbols.TryGetValue(symName, symbolFiles) and (Length(symbolFiles) > 0) then + if (Symbols.TryGetValue(symName, symbolFiles)) and (Length(symbolFiles) > 0) then Result := symbolFiles[0] else Result := Default(TDataFile); @@ -287,7 +305,7 @@ begin foundIndex := -1; for var i := 0 to High(allSymbolFiles) do begin - if allSymbolFiles[i].GetBaseFileName = Curr.GetBaseFileName then + if (allSymbolFiles[i].GetBaseFileName = Curr.GetBaseFileName) then begin foundIndex := i; break; @@ -314,39 +332,37 @@ begin ext := TPath.GetExtension(fileNameNoPath); baseName := TPath.GetFileNameWithoutExtension(fileNameNoPath); - // Format: SYMBOL_YYYY_MM parts := baseName.Split(['_']); - if Length(parts) < 3 then + if (Length(parts) < 3) then exit; - // Use reverse indexing for year and month if (not TryStrToInt(parts[High(parts)], month)) or (not TryStrToInt(parts[High(parts) - 1], year)) then exit; - symbol := string.Join('_', Copy(parts, 0, Length(parts) - 2)); - Result := TDataFile.Create(TPath.GetDirectoryName(FileName), symbol, ext, year, month, baseName); + symbol := string.Join('_', Copy(parts, 0, Length(parts) - 2)) + ext; + Result := TDataFile.Create(TPath.GetDirectoryName(FileName), symbol, year, month, baseName); end; -function TDataServer.LoadDataFile(const DataFile: TDataFile): TFuture>>>; +function TDataServer.LoadDataFile(const DataFile: TDataFile): TFuture>; begin - if DataFile.IsValid then + if (DataFile.IsValid) then Result := FCachedFiles.GetOrAdd( DataFile.GetFullFileName, - function(const Filename: String): TFuture>>> begin Result := DoLoad(Filename, FLoadGate); end + function(const Filename: String): TFuture> begin Result := DoLoad(Filename, FLoadGate); end ) else - Result := TFuture>>>.Construct(TArray>>(nil)); + Result := TFuture>.Construct(TArray(nil)); end; -function TDataServer.DoLoad(const FileName: string; const Gate: TLatch): TFuture>>>; +function TDataServer.DoLoad(const FileName: string; const Gate: TLatch): TFuture>; begin var capFileName := FileName; Result := TFuture .Construct(Gate.Enqueue, function: TBytes begin Result := TFile.ReadAllBytes(capFileName); end) - .Chain>>>( - function(Bytes: TBytes): TArray>> + .Chain>( + function(Bytes: TBytes): TArray var compStream: TBytesStream; begin @@ -359,7 +375,7 @@ begin end); end; -function TDataServer.ReadCompressedData(const InputStream: TStream): TArray>>; +function TDataServer.ReadCompressedData(const InputStream: TStream): TArray; var zip: TZipFile; decompStream: TStream; @@ -375,14 +391,14 @@ begin // Search for the .bin entry created by C# for var i := 0 to zip.FileCount - 1 do begin - if zip.FileNames[i].EndsWith('.bin', true) then + if (zip.FileNames[i].EndsWith('.bin', true)) then begin entryIdx := i; break; end; end; - if entryIdx = -1 then + if (entryIdx = -1) then exit; zip.Read(entryIdx, decompStream, header); @@ -396,37 +412,54 @@ begin end; end; -function TDataServer.ReadUncompressedData(const InputStream: TStream): TArray>>; +function TDataServer.ReadUncompressedData(const InputStream: TStream): TArray; var totalRecords, chunkCount, i, j, recsInChunk: Integer; - rawBuffer: TArray; + inputRecordSize, outputPointSize: Integer; + chunkBuffer: TArray>; + rawBufferChunk: TArray; // Only buffer one chunk at a time begin Result := nil; - totalRecords := InputStream.Size div sizeof(R); - if totalRecords = 0 then - exit; + inputRecordSize := SizeOf(R); + outputPointSize := SizeOf(TDataPoint); - // Fast block read into memory - SetLength(rawBuffer, totalRecords); - InputStream.Position := 0; - InputStream.ReadBuffer(rawBuffer[0], InputStream.Size); + // Ensure we are working with complete records + if (InputStream.Size mod inputRecordSize <> 0) then + raise Exception.Create('Stream size is not a multiple of record size.'); + + totalRecords := InputStream.Size div inputRecordSize; + if (totalRecords = 0) then + exit; chunkCount := (totalRecords + ChunkSize - 1) div ChunkSize; SetLength(Result, chunkCount); + SetLength(rawBufferChunk, ChunkSize); // Reuse buffer + + InputStream.Position := 0; for i := 0 to chunkCount - 1 do begin recsInChunk := Min(ChunkSize, totalRecords - (i * ChunkSize)); - SetLength(Result[i], recsInChunk); + + // Read only needed bytes for this chunk + InputStream.ReadBuffer(rawBufferChunk[0], recsInChunk * inputRecordSize); + + // Temp buffer for mapped objects + SetLength(chunkBuffer, recsInChunk); + for j := 0 to recsInChunk - 1 do begin - Result[i][j] := MapRecord(rawBuffer[i * ChunkSize + j]); + chunkBuffer[j] := MapRecord(rawBufferChunk[j]); end; + + // Serialize mapped data to TBytes + SetLength(Result[i], recsInChunk * outputPointSize); + Move(chunkBuffer[0], Result[i][0], recsInChunk * outputPointSize); end; end; function TDataServer.ProcessChunks( - const DataChunks: TArray>>; + const MappedBlobs: TArray; Terminated: TState; Processor: IConsumer>> ): TState; @@ -434,16 +467,28 @@ begin var callProcess := function(Idx: Integer): TFunc begin - var cChunk := DataChunks[Idx]; + // Capture the blob + var cBlob := MappedBlobs[Idx]; Result := function: TState + var + dataChunk: TArray>; + count: Integer; begin if not Terminated.IsSet then - Result := Processor.Consume(cChunk); + begin + // Quick deserialize (Mem Copy) + count := Length(cBlob) div SizeOf(TDataPoint); + SetLength(dataChunk, count); + if count > 0 then + Move(cBlob[0], dataChunk[0], Length(cBlob)); + + Result := Processor.Consume(dataChunk); + end; end; end; - for var i := 0 to High(DataChunks) do + for var i := 0 to High(MappedBlobs) do begin if not Terminated.IsSet then begin @@ -455,7 +500,7 @@ end; function TDataServer.ProcessFile( FileInfo: TDataFile; - const DataFile: TFuture>>>; + const DataFile: TFuture>; const Terminated: TState; Processor: IConsumer>> ): TState; @@ -469,9 +514,9 @@ begin Result := DataFile.Chain( - function(const DataChunks: TArray>>): TState + function(const MappedBlobs: TArray): TState begin - var done := ProcessChunks(DataChunks, cTerminated, Processor); + var done := ProcessChunks(MappedBlobs, cTerminated, Processor); Result := TaskManager .RunTask(done, function: TState begin Result := ProcessFile(nextFileInfo, nextFile, cTerminated, Processor); end);