From 342eb07c42e11cab6d54332c32cc2df8d10b7614 Mon Sep 17 00:00:00 2001 From: Michael Schimmel Date: Wed, 16 Jul 2025 14:12:07 +0200 Subject: [PATCH] Fixed concurrent processing --- AuraTrader/MainForm.pas | 44 ++++++++++-------- Src/Myc.Fmx.Chart.Series.pas | 38 ++------------- Src/Myc.Fmx.Chart.pas | 51 +++++++++----------- Src/Myc.TaskManager.pas | 20 +++++++- Src/Myc.Trade.DataPoint.Impl.pas | 80 +++++++++++++++++++++++++++----- Src/Myc.Trade.DataPoint.pas | 27 +++++++++++ Src/Myc.Trade.DataStream.pas | 36 ++++++++++---- Src/Myc.Trade.Types.pas | 1 + 8 files changed, 192 insertions(+), 105 deletions(-) diff --git a/AuraTrader/MainForm.pas b/AuraTrader/MainForm.pas index 5997f3d..069099a 100644 --- a/AuraTrader/MainForm.pas +++ b/AuraTrader/MainForm.pas @@ -108,6 +108,7 @@ type TEquitySum = class(TMycConverter) private FEquity: Double; + FInit: Boolean; protected function ProcessData(const Value: Double): TState; override; public @@ -360,12 +361,12 @@ begin var Ohlc := OhlcPoint.Field('Data'); var Closes := Ohlc.Field('Close'); - var Hull := Closes.Chain(TIndicators.CreateHMA(150)); + var Hull := Closes.MakeParallel.Chain(TIndicators.CreateHMA(150)); var Sma := Closes.Chain(TIndicators.CreateSMA(50)); var Ema := Closes.Chain(TIndicators.CreateEMA(21)); - var Boli := Closes.Chain(TIndicators.CreateBollingerBands(20, 2.0)); - var Rsi := Closes.Chain(TIndicators.CreateRSI(14)); - var Macd := Closes.Chain(TIndicators.CreateMACD(12, 26, 9)); + var Boli := Closes.MakeParallel.Chain(TIndicators.CreateBollingerBands(20, 2.0)); + var Rsi := Closes.MakeParallel.Chain(TIndicators.CreateRSI(14)); + var Macd := Closes.MakeParallel.Chain(TIndicators.CreateMACD(12, 26, 9)); var Stoch := Ohlc.Chain(TIndicators.CreateStochastic(14, 3)); chart.SetXAxisSeries(timeframe, Timestamps.Sender); @@ -429,7 +430,7 @@ type pnl: Double; end; begin - var timeframe := TTimeframe.M15; + var timeframe := TTimeframe.D; var ticker := TConverter.CreateTicker>; @@ -440,37 +441,35 @@ begin ) ); - var OhlcPoint := lastPrice.Chain>(TConverter.CreateAggregation(timeframe)); + var OhlcPoint := lastPrice.Chain>(TConverter.CreateAggregation(timeframe)).MakeParallel; var Ohlc := TConverter.CreateSequence(2, OhlcPoint.Field('Data').Sender); var Closes := Ohlc[0].Field('Close'); - var Hull := Closes.Chain(TIndicators.CreateHMA(250)); - var Sma := Closes.Chain(TIndicators.CreateSMA(200)); - - var HullSeries := TConverter.CreateEndpoint(Hull.Sender, 5); - var SmaSeries := TConverter.CreateEndpoint(Sma.Sender, 5); + var Hull := Closes.Chain(TIndicators.CreateHMA(250)).MakeParallel; + var Sma := Closes.Chain(TIndicators.CreateSMA(200)).MakeParallel; var Lowest: Double := Double.MaxValue; var Highest: Double := Double.MinValue; + var ATR := Ohlc[0].Chain(TIndicators.CreateATR(15)); + var ATRSeries := TConverter.CreateEndpoint(ATR.Sender, 5); + + var HullSeries := TConverter.CreateEndpoint(Hull.Sender, 5); + var SmaSeries := TConverter.CreateEndpoint(Sma.Sender, 5); + + // next stage + var curr: TSignal; curr.SL := Double.NaN; curr.Entry := Double.NaN; - var ATR := Ohlc[0].Chain(TIndicators.CreateATR(50)); - var ATRSeries := TConverter.CreateEndpoint(ATR.Sender, 5); - - // next stage - var Signal := Ohlc[1] .Chain( function(const Ohlc: TOhlcItem): TSignal begin - var pnl: Double := 0; - if Ohlc.Low < Lowest then Lowest := Ohlc.Low; if Ohlc.High > Highest then @@ -478,7 +477,7 @@ begin Result := curr; Result.Sig := 0; - pnl := NaN; + var pnl: double := NaN; if (HullSeries.Value[0] < SmaSeries.Value[0]) and (HullSeries.Value[1] >= SmaSeries.Value[1]) then begin @@ -600,10 +599,17 @@ constructor TEquitySum.Create(AEquity: Double); begin inherited Create; FEquity := AEquity; + FInit := false; end; function TEquitySum.ProcessData(const Value: Double): TState; begin + if not FInit then + begin + FInit := true; + Broadcast(FEquity); + end; + if not IsNan(Value) then begin FEquity := FEquity + Value; diff --git a/Src/Myc.Fmx.Chart.Series.pas b/Src/Myc.Fmx.Chart.Series.pas index cd428a3..798e690 100644 --- a/Src/Myc.Fmx.Chart.Series.pas +++ b/Src/Myc.Fmx.Chart.Series.pas @@ -15,27 +15,13 @@ uses Myc.Fmx.Chart; type - TChartSeriesReceiver = class(TMycProcessor) - strict private - FCurrData: TSeries; - FLookback: TMutable; - private - FData: TWriteable>; - protected - function ProcessData(const Value: T): TState; override; - public - constructor Create(const ALookback: TMutable); - property Data: TWriteable> read FData; - end; - TChartSeriesProcessor = class(TMycChart.TSeries) strict private FDataSeries: TSeries; private FData: TSeries; FDataProvider: TDataProvider; - FReceiver: TChartSeriesReceiver; - FReceiverTag: TDataProvider.TTag; + FReceiver: TMutable>; protected function GetCount: Int64; override; function GetTotalCount: Int64; override; @@ -140,22 +126,6 @@ uses System.Math, FMX.Types; -{ TChartSeriesReceiver } - -constructor TChartSeriesReceiver.Create(const ALookback: TMutable); -begin - inherited Create; - FLookback := ALookback; - FData := TWriteable>.CreateWriteable(FCurrData).Protect; -end; - -function TChartSeriesReceiver.ProcessData(const Value: T): TState; -begin - Result := TState.Null; - FCurrData := FCurrData.Add(Value, FLookback.Value); - FData.Value := FCurrData; -end; - { TChartSeriesProcessor } constructor TChartSeriesProcessor.Create(const ADataProvider: TDataProvider; const ALookback: TMutable); @@ -164,13 +134,11 @@ begin FDataProvider := ADataProvider; FData := FDataSeries; - FReceiver := TChartSeriesReceiver.Create(ALookback); - FReceiverTag := FDataProvider.Link(FReceiver); + FReceiver := TConverter.CreateEndpoint(FDataProvider, ALookback.Value); end; destructor TChartSeriesProcessor.Destroy; begin - FDataProvider.Unlink(FReceiverTag); inherited; end; @@ -186,7 +154,7 @@ end; procedure TChartSeriesProcessor.Update; begin - FData := FReceiver.Data.Value; + FData := FReceiver.Value; end; { TChartLineLayer } diff --git a/Src/Myc.Fmx.Chart.pas b/Src/Myc.Fmx.Chart.pas index eb69e03..b0a27bd 100644 --- a/Src/Myc.Fmx.Chart.pas +++ b/Src/Myc.Fmx.Chart.pas @@ -64,18 +64,19 @@ type constructor Create(AParent: TPanel); end; + TDragPoint = record + Point: TPointF; + Idx: Int64; + BarX: Single; + end; + TAxisLayer = class abstract(TLayer) private FOwner: TMycChart; protected function GetOwner: TMycChart; override; final; // Paint crosshair and other axis related overlays - procedure Paint( - const Canvas: TCanvas; - const Viewport: TRectF; - ViewStartIndex, ViewCount: Int64; - IdxAtMousePos: Integer - ); virtual; abstract; + procedure Paint(const Canvas: TCanvas; const Viewport: TRectF; const MousePos: TDragPoint); virtual; abstract; public constructor Create(AOwner: TMycChart); end; @@ -83,12 +84,7 @@ type TXAxisLayer = class(TAxisLayer) protected // Paint crosshair and time caption - procedure Paint( - const Canvas: TCanvas; - const Viewport: TRectF; - ViewStartIndex, ViewCount: Int64; - IdxAtMousePos: Integer - ); override; + procedure Paint(const Canvas: TCanvas; const Viewport: TRectF; const MousePos: TDragPoint); override; // This function delivers the text for the caption. function GetCaption(Idx: Int64): String; virtual; abstract; end; @@ -124,11 +120,6 @@ type property Weight: Single read FWeight write SetWeight; end; - TDragPoint = record - Point: TPointF; - Idx: Int64; - end; - private FPanelList: TObjectList; FXAxisSeries: TMycChart.TXAxisLayer; @@ -175,6 +166,8 @@ type property Panels[Index: Integer]: TPanel read GetPanel; default; // The rectangle occupied by the jump-to-latest button, used for hit testing. property JumpButtonRect: TRectF read FJumpButtonRect; + property ViewCount: Int64 read FViewCount; + property ViewStartIndex: Int64 read FViewStartIndex; end; implementation @@ -232,11 +225,16 @@ end; function TMycChart.CreateDragPoint(X, Y: Single): TDragPoint; begin + var rect := LocalRect; + Result.Point.X := X; Result.Point.Y := Y; - Result.Idx := Round((LocalRect.Right - X) * (FViewCount - 1) / LocalRect.Width); + Result.Idx := Round((rect.Right - X) * (FViewCount - 1) / rect.Width); if (Result.Idx < 0) or (Result.Idx >= FXAxisSeries.Series.Count) then Result.Idx := -1; + Result.BarX := NaN; + if Result.Idx >= 0 then + Result.BarX := rect.Right - (Result.Idx * (FViewCount - 1)) / (FViewCount - 1) * (rect.Width / (FViewCount - 1)); end; function TMycChart.GetPanel(Index: Integer): TPanel; @@ -356,7 +354,7 @@ begin begin // Do not draw crosshair if mouse is over the jump button if not FJumpButtonRect.Contains(FMousePos.Point) then - FXAxisSeries.Paint(Self.Canvas, rect, FViewStartIndex, FViewCount, FViewStartIndex + FMousePos.Idx); + FXAxisSeries.Paint(Self.Canvas, rect, FMousePos); end; // --- Draw the "jump to latest" button (code unchanged) --- @@ -839,12 +837,7 @@ end; { TMycChart.TXAxisLayer } -procedure TMycChart.TXAxisLayer.Paint( - const Canvas: TCanvas; - const Viewport: TRectF; - ViewStartIndex, ViewCount: Int64; - IdxAtMousePos: Integer -); +procedure TMycChart.TXAxisLayer.Paint(const Canvas: TCanvas; const Viewport: TRectF; const MousePos: TDragPoint); var caption: string; captionRect: TRectF; @@ -853,11 +846,11 @@ var textSize: TRectF; begin // Do not draw if no data or view is invalid - if (not Assigned(Series)) or (Series.Count = 0) or (ViewCount <= 1) then + if (not Assigned(Series)) or (Series.Count = 0) then exit; - // Snap X to the candle's center - candleX := Viewport.Right - ((IdxAtMousePos - ViewStartIndex) * (ViewCount - 1)) / (ViewCount - 1) * (Viewport.Width / (ViewCount - 1)); + // 1. Snap X to the candle's center + candleX := MousePos.BarX; // 2. Draw the vertical line at the snapped position Canvas.Stroke.Kind := TBrushKind.Solid; @@ -866,7 +859,7 @@ begin Canvas.DrawLine(PointF(candleX, Viewport.Top), PointF(candleX, Viewport.Bottom), 1.0); // 3. Get the caption for the index - caption := GetCaption(IdxAtMousePos); + caption := GetCaption(FOwner.ViewStartIndex + MousePos.Idx); if caption.IsEmpty then exit; diff --git a/Src/Myc.TaskManager.pas b/Src/Myc.TaskManager.pas index b8dc021..cb05cec 100644 --- a/Src/Myc.TaskManager.pas +++ b/Src/Myc.TaskManager.pas @@ -29,7 +29,10 @@ type // Run a task when the gate is opened. // Returns a State to await thread completion. - function RunTask(const Gate: TState; const Proc: TFunc): TState; + function RunTask(const Proc: TFunc): TState; overload; + function RunTask(const Gate: TState; const Proc: TFunc): TState; overload; + procedure RunTask(const Gate: TState; const Proc: TProc); overload; + procedure RunTask(const Proc: TProc); overload; class function RunSequence(const Gate: TState; First, Count: Integer; const Proc: TFunc): TState; static; @@ -181,6 +184,21 @@ begin ); end; +function TTaskManager.RunTask(const Proc: TFunc): TState; +begin + Result := RunTask(TState.Null, Proc); +end; + +procedure TTaskManager.RunTask(const Proc: TProc); +begin + RunTask(TState.Null, Proc); +end; + +procedure TTaskManager.RunTask(const Gate: TState; const Proc: TProc); +begin + FTaskManager.RunTask(Gate, Proc); +end; + procedure TTaskManager.WaitFor(const State: TState); begin FTaskManager.WaitFor(State); diff --git a/Src/Myc.Trade.DataPoint.Impl.pas b/Src/Myc.Trade.DataPoint.Impl.pas index c567d84..c1107e4 100644 --- a/Src/Myc.Trade.DataPoint.Impl.pas +++ b/Src/Myc.Trade.DataPoint.Impl.pas @@ -4,6 +4,7 @@ interface uses System.SysUtils, + System.SyncObjs, Myc.Signals, Myc.Mutable, Myc.Core.Notifier, @@ -15,6 +16,7 @@ type // Abstract base class for data consumers. TMycProcessor = class abstract(TInterfacedObject, IMycProcessor) protected + function Execute(const Value: T; const Proc: TConstFunc): TState; function ProcessData(const Value: T): TState; virtual; abstract; end; @@ -86,7 +88,16 @@ type constructor Create(const AFunc: TConstFunc); end; - // A converter specialized for calculating indicators. + TMycGenericParallelConverter = class(TMycConverter) + private + FFunc: TConstFunc; + FQueue: TState; + protected + function ProcessData(const Value: S): TState; override; + public + constructor Create(const AFunc: TConstFunc); + end; + TMycIdentityConverter = class(TMycConverter) protected function ProcessData(const Value: T): TState; override; final; @@ -126,14 +137,12 @@ type // A generic processor implementation that uses a function reference. TMycGenericProcessor = class(TMycProcessor) - type - TProc = reference to function(const Value: T): TState; private - FProc: TProc; + FProc: TConstFunc; protected function ProcessData(const Value: T): TState; override; final; public - constructor Create(const AProc: TProc); + constructor Create(const AProc: TConstFunc); end; // A processor implementation that is owned by a controller. @@ -156,6 +165,7 @@ type FLookback: Int64; FData: TSeries; FChanged: TEvent; + FLock: TLightweightMREW; function GetChanged: TSignal; function GetValue: TSeries; function ProcessData(const Value: T): TState; @@ -178,6 +188,13 @@ type property Timeframe: TTimeframe read GetTimeframe; end; + TMycParallelConverter = class(TMycConverter) + private + FQueue: TState; + protected + function ProcessData(const Value: T): TState; override; final; + end; + implementation uses @@ -187,6 +204,14 @@ uses System.Math, Myc.TaskManager; +{ TMycProcessor } + +function TMycProcessor.Execute(const Value: T; const Proc: TConstFunc): TState; +begin + var cValue := Value; + Result := TaskManager.RunTask(function: TState begin Result := Proc(cValue); end); +end; + { TMycContainedDataProvider } constructor TMycContainedDataProvider.Create(const Controller: IInterface); @@ -314,7 +339,7 @@ end; function TMycIndicator.ProcessData(const Value: S): TState; begin - Result := Broadcast(Calculate(Value)); + Result := Execute(Value, function(const Value: S): TState begin Result := Broadcast(Calculate(Value)); end); end; { TMycDataCounter } @@ -385,7 +410,7 @@ end; { TMycGenericProcessor } -constructor TMycGenericProcessor.Create(const AProc: TProc); +constructor TMycGenericProcessor.Create(const AProc: TConstFunc); begin inherited Create; FProc := AProc; @@ -419,6 +444,7 @@ begin FProcessor := TMycContainedProcessor.Create(Self, ProcessData); FTag := FDataProvider.Link(FProcessor); + end; destructor TMycDataEndpoint.Destroy; @@ -435,15 +461,24 @@ end; function TMycDataEndpoint.GetValue: TSeries; begin - Result := FData; + FLock.BeginRead; + try + Result := FData; + finally + FLock.EndRead; + end; end; function TMycDataEndpoint.ProcessData(const Value: T): TState; begin Result := TState.Null; - - // Add new data point, respecting the lookback period. - FData := FData.Add(Value, FLookback); + FLock.BeginWrite; + try + // Add new data point, respecting the lookback period. + FData := FData.Add(Value, FLookback); + finally + FLock.EndWrite; + end; FChanged.Notify; end; @@ -599,4 +634,27 @@ begin Result := Broadcast(Value); end; +{ TMycGenericParallelConverter } + +constructor TMycGenericParallelConverter.Create(const AFunc: TConstFunc); +begin + inherited Create; + FFunc := AFunc; +end; + +function TMycGenericParallelConverter.ProcessData(const Value: S): TState; +begin + var cValue := Value; + Result := TaskManager.RunTask(FQueue, function: TState begin Result := Broadcast(FFunc(cValue)); end); + + FQueue := Result; +end; + +function TMycParallelConverter.ProcessData(const Value: T): TState; +begin + var cValue := Value; + Result := TaskManager.RunTask(FQueue, function: TState begin Result := Broadcast(cValue); end); + FQueue := Result; +end; + end. diff --git a/Src/Myc.Trade.DataPoint.pas b/Src/Myc.Trade.DataPoint.pas index 8725cba..3e0c07b 100644 --- a/Src/Myc.Trade.DataPoint.pas +++ b/Src/Myc.Trade.DataPoint.pas @@ -82,9 +82,13 @@ type class operator Implicit(const A: TConverter): IConverter; overload; class function CreateGeneric(const Func: TConstFunc): TConverter; static; + class function CreateParallel(const Func: TConstFunc): TConverter; static; function Chain(const Next: TConverter): TConverter; overload; inline; function Chain(const Func: TConstFunc): TConverter; overload; inline; + function ChainParallel(const Func: TConstFunc): TConverter; overload; inline; + + function MakeParallel: TConverter; overload; inline; // Extracts the field of a record by it's name (using RTTI). function Field(const FieldName: String): TConverter; overload; inline; @@ -112,6 +116,8 @@ type class function CreateDataPointConverter(const Func: TConstFunc): TConverter, TDataPoint>; static; class function CreateSequence(Count: Integer; const Parent: TDataProvider): TArray>; overload; static; + + class function Parallel(Parent: TDataProvider): TConverter; static; end; implementation @@ -184,11 +190,21 @@ begin Result := Chain(TMycGenericConverter.Create(Func)); end; +function TConverter.ChainParallel(const Func: TConstFunc): TConverter; +begin + Result := Chain(TMycGenericParallelConverter.Create(Func)); +end; + class function TConverter.CreateGeneric(const Func: TConstFunc): TConverter; begin Result := TMycGenericConverter.Create(Func); end; +class function TConverter.CreateParallel(const Func: TConstFunc): TConverter; +begin + Result := TMycGenericParallelConverter.Create(Func); +end; + function TConverter.Field(const FieldName: String): TConverter; begin Result := Chain(TConverter.CreateRecordField(FieldName)); @@ -199,6 +215,11 @@ begin Result := FConverter.Sender; end; +function TConverter.MakeParallel: TConverter; +begin + Result := Chain(TMycGenericParallelConverter.Create(function(const Value: T): T begin Result := Value; end)); +end; + function TConverter.Sequence(Count: Integer): IMycDataSequence; begin Result := TMycSequence.Create(Count); @@ -288,4 +309,10 @@ begin Parent.Link(seq); end; +class function TConverter.Parallel(Parent: TDataProvider): TConverter; +begin + Result := TMycParallelConverter.Create; + Parent.Link(Result); +end; + end. diff --git a/Src/Myc.Trade.DataStream.pas b/Src/Myc.Trade.DataStream.pas index 635361b..d053857 100644 --- a/Src/Myc.Trade.DataStream.pas +++ b/Src/Myc.Trade.DataStream.pas @@ -60,6 +60,12 @@ type // Used by cache to actually load an uncached file. function DoLoad(const FileName: string): TFuture>>>; + function ProcessChunks( + const DataChunks: TArray>>; + const Terminated: TState; + Processor: IMycProcessor>> + ): TState; + function ProcessFile( FileInfo: TAuraDataFile; const DataFile: TFuture>>>; @@ -354,6 +360,23 @@ begin Result := FCachedFiles.GetOrAdd(DataFile.GetFullFileName); end; +function TAuraDataServer.ProcessChunks( + const DataChunks: TArray>>; + const Terminated: TState; + Processor: IMycProcessor>> +): TState; +begin + var done := TLatch.CreateLatch(Length(DataChunks)); + + for var i := 0 to High(DataChunks) do + if not Terminated.IsSet then + Processor.ProcessData(DataChunks[i]).Signal.Subscribe(done) + else + done.Notify; + + Result := done.State; +end; + function TAuraDataServer.ProcessFile( FileInfo: TAuraDataFile; const DataFile: TFuture>>>; @@ -374,19 +397,12 @@ begin function(const DataChunks: TArray>>): TState begin // Process each chunk, checking for termination between chunks. - var done := TLatch.CreateLatch(Length(DataChunks)); - for var i := 0 to High(DataChunks) do - if not cTerminated.IsSet then - Processor.ProcessData(DataChunks[i]).Signal.Subscribe(done) - else - done.Notify; + var done := ProcessChunks(DataChunks, Terminated, Processor); // Move to next file (which is currently preloading) once all Processors are done Result := - TaskManager.RunTask( - done.State, - function: TState begin Result := ProcessFile(nextFileInfo, nextFile, cTerminated, Processor); end - ) + TaskManager + .RunTask(done, function: TState begin Result := ProcessFile(nextFileInfo, nextFile, cTerminated, Processor); end); end ); end; diff --git a/Src/Myc.Trade.Types.pas b/Src/Myc.Trade.Types.pas index 551f2d4..d7012b0 100644 --- a/Src/Myc.Trade.Types.pas +++ b/Src/Myc.Trade.Types.pas @@ -29,6 +29,7 @@ type end; TConstFunc = reference to function(const Value: S): T; + TConstProc = reference to procedure(const Value: T); TConstFuncPredicate = reference to function(const Value: S; out Res: T): Boolean; implementation