From 13c41d01b54c8fa24c793781e1da9430b824f996 Mon Sep 17 00:00:00 2001 From: Michael Schimmel Date: Tue, 24 Jun 2025 20:17:04 +0200 Subject: [PATCH] TaskManager.RunTask() & TFuture.Chain(Proc) --- Src/Myc.Core.Futures.pas | 35 ++++++------- Src/Myc.Core.Tasks.pas | 42 ++++++---------- Src/Myc.Futures.pas | 16 +++--- Src/Myc.Signals.pas | 25 ++++++---- Src/Myc.TaskManager.pas | 65 +++++++++++++++++++----- Src/Myc.Trade.DataStream.pas | 96 +++++------------------------------- Test/TestCoreFutures.pas | 4 +- Test/TestTasks.pas | 2 +- 8 files changed, 120 insertions(+), 165 deletions(-) diff --git a/Src/Myc.Core.Futures.pas b/Src/Myc.Core.Futures.pas index 27ca3c2..26801b1 100644 --- a/Src/Myc.Core.Futures.pas +++ b/Src/Myc.Core.Futures.pas @@ -25,14 +25,13 @@ type TMycGateFuncFuture = class(TMycFuture) private - FInit: TSignal.TSubscription; FDone: TLatch.ILatch; FResult: T; protected function GetValue: T; override; function GetDone: TState; override; public - constructor Create(const ATaskManager: IMycTaskManager; const AGate: TState.IState; AProc: TFunc); + constructor Create(const ATaskManager: TTaskManager; const AGate: TState.IState; AProc: TFunc); destructor Destroy; override; end; @@ -71,7 +70,7 @@ end; { TMycGateFuncFuture } -constructor TMycGateFuncFuture.Create(const ATaskManager: IMycTaskManager; const AGate: TState.IState; AProc: TFunc); +constructor TMycGateFuncFuture.Create(const ATaskManager: TTaskManager; const AGate: TState.IState; AProc: TFunc); begin inherited Create; @@ -79,28 +78,26 @@ begin // Subscribe the job execution to AGate. // The job will run when AGate notifies the subscriber returned by Run. - FInit := - ATaskManager.CreateTask( - AGate, - procedure - begin + ATaskManager.RunTask( + AGate, + procedure + begin + try try - try - Self.FResult := AProc(); - except - Self.FResult := Default(T); // Set result to Default(T) on error - raise; // Re-raise for TaskFactory to handle - end; - finally - Self.FDone.Notify; // Signal that this future is done (successfully or with error) + Self.FResult := AProc(); + except + Self.FResult := Default(T); // Set result to Default(T) on error + raise; // Re-raise for TaskFactory to handle end; - end - ); + finally + Self.FDone.Notify; // Signal that this future is done (successfully or with error) + end; + end + ); end; destructor TMycGateFuncFuture.Destroy; begin - FInit.Unsubscribe; inherited Destroy; end; diff --git a/Src/Myc.Core.Tasks.pas b/Src/Myc.Core.Tasks.pas index 6feae11..c198f45 100644 --- a/Src/Myc.Core.Tasks.pas +++ b/Src/Myc.Core.Tasks.pas @@ -11,7 +11,7 @@ uses Myc.Signals; type - IMycTaskFactory = interface(IMycTaskManager) + IMycThreadPool = interface(TTaskManager.IMycTaskManager) {$region 'property access'} // Retrieves the number of worker threads. function GetThreadCount: Integer; @@ -49,7 +49,7 @@ type property ThreadCount: Integer read GetThreadCount; end; - TMycTaskFactory = class(TInterfacedObject, IMycTaskManager, IMycTaskFactory) + TMycTaskFactory = class(TInterfacedObject, TTaskManager.IMycTaskManager, IMycThreadPool) type ETaskException = class(Exception) end; @@ -65,9 +65,7 @@ type FWorkStack: TMycAtomicStack; // Stack of pending jobs FWorkThreads: TArray; // Array of worker threads procedure WorkerThread; - class var - FTaskManagerLock: Integer; - + class constructor CreateClass; protected procedure ExecuteJob; function GetThreadCount: Integer; @@ -81,16 +79,12 @@ type procedure HandleException; procedure EnqueueJob(const Job: TProc); function Run(Job: TProc): TSignal.ISubscriber; - function CreateTask(const Gate: TState; const Proc: TProc): TSignal.TSubscription; - procedure WaitFor(const State: TState.IState); + procedure RunTask(const Gate: TState; const Proc: TProc); + procedure WaitFor(const State: TState); procedure Teardown; function InMainThread: Boolean; function InWorkerThread: Boolean; - // Initialize global TaskManager - class procedure AquireTaskManager; - class procedure ReleaseTaskManager; - property ThreadsRunning: Integer read FThreadsRunning; property WorkThreads: TArray read FWorkThreads; end; @@ -147,6 +141,10 @@ begin FWaitSemaphores.Push(TSemaphore.Create(nil, 0, 1, '')); end; +class constructor TMycTaskFactory.CreateClass; +begin +end; + destructor TMycTaskFactory.Destroy; begin Teardown; // Perform cleanup and stop threads @@ -188,7 +186,7 @@ begin Result.NameThreadForDebugging(DbgName); // Set thread name for debugging end; -function TMycTaskFactory.CreateTask(const Gate: TState; const Proc: TProc): TSignal.TSubscription; +procedure TMycTaskFactory.RunTask(const Gate: TState; const Proc: TProc); begin if Gate.IsSet then begin @@ -197,7 +195,7 @@ begin else begin // The job will run when Gate notifies the subscriber returned by Run. - Result := Gate.Signal.Subscribe(Run(Proc)); + Gate.Signal.Subscribe(Run(Proc)); end; end; @@ -289,20 +287,6 @@ begin exit(false); // Not found in worker thread list end; -class procedure TMycTaskFactory.AquireTaskManager; -begin - if not Assigned(TaskManager) then - TaskManager := TMycTaskFactory.Create; - inc(FTaskManagerLock); -end; - -class procedure TMycTaskFactory.ReleaseTaskManager; -begin - dec(FTaskManagerLock); - if FTaskManagerLock = 0 then - TaskManager := nil; -end; - function TMycTaskFactory.Run(Job: TProc): TSignal.ISubscriber; begin if FTerminated <> 0 then @@ -345,7 +329,7 @@ begin end; end; -procedure TMycTaskFactory.WaitFor(const State: TState.IState); +procedure TMycTaskFactory.WaitFor(const State: TState); var lock: TSemaphore; // This is System.SyncObjs.TSemaphore begin @@ -431,4 +415,6 @@ end; initialization IsMultiThread := true; + TaskManager := TMycTaskFactory.Create; + end. diff --git a/Src/Myc.Futures.pas b/Src/Myc.Futures.pas index 67f5b66..7ac3bb2 100644 --- a/Src/Myc.Futures.pas +++ b/Src/Myc.Futures.pas @@ -23,6 +23,7 @@ type end; TFuncConst = reference to function(const Arg1: S): TResult; + TProcConst = reference to function(const Arg1: S): TResult; {$REGION 'private'} strict private @@ -30,8 +31,6 @@ type FNull: IFuture; class constructor CreateClass; - class destructor DestroyClass; - private FFuture: IFuture; function GetDone: TState.IState; inline; @@ -53,6 +52,7 @@ type function Chain(const Proc: TFunc): TFuture; overload; function Chain(const Proc: TFuncConst): TFuture; overload; + function Chain(const Proc: TProcConst): TState; overload; function WaitFor: T; property Done: TState.IState read GetDone; @@ -62,8 +62,7 @@ type implementation uses - Myc.Core.Futures, - Myc.Core.Tasks; + Myc.Core.Futures; constructor TFuture.Create(const AFuture: IFuture); begin @@ -74,13 +73,16 @@ end; class constructor TFuture.CreateClass; begin - TMycTaskFactory.AquireTaskManager; FNull := TMycNullFuture.Create; end; -class destructor TFuture.DestroyClass; +function TFuture.Chain(const Proc: TProcConst): TState; begin - TMycTaskFactory.ReleaseTaskManager; + var Done := TLatch.CreateLatch(1); + Result := Done.State; + + var future := FFuture; + TaskManager.RunTask(future.Done, procedure begin Proc(future.Value).Subscribe(Done); end); end; function TFuture.Chain(const Proc: TFunc): TFuture; diff --git a/Src/Myc.Signals.pas b/Src/Myc.Signals.pas index 1b2e2e2..0361807 100644 --- a/Src/Myc.Signals.pas +++ b/Src/Myc.Signals.pas @@ -233,8 +233,9 @@ end; constructor TState.Create(const AState: IState); begin - if Assigned(AState) then - FState := AState; + FState := AState; + if not Assigned(FState) then + FState := Null; end; class constructor TState.ClassCreate; @@ -304,8 +305,9 @@ end; constructor TEvent.Create(const AEvent: IEvent); begin - if Assigned(AEvent) then - FEvent := AEvent; + FEvent := AEvent; + if not Assigned(FEvent) then + FEvent := Null; end; class constructor TEvent.ClassCreate; @@ -357,8 +359,9 @@ end; constructor TFlag.Create(const AFlag: IFlag); begin - if Assigned(AFlag) then - FFlag := AFlag; + FFlag := AFlag; + if not Assigned(FFlag) then + FFlag := Null; end; class function TFlag.CreateFlag(Init: Boolean = false): TFlag; @@ -412,8 +415,9 @@ end; constructor TLatch.Create(const ALatch: TLatch.ILatch); begin - if Assigned(ALatch) then - FLatch := ALatch; + FLatch := ALatch; + if not Assigned(FLatch) then + FLatch := Null; end; class function TLatch.CreateLatch(Count: Integer): ILatch; @@ -500,8 +504,9 @@ end; constructor TSignal.Create(const ASignal: ISignal); begin - if Assigned(FSignal) then - FSignal := ASignal; + FSignal := ASignal; + if not Assigned(FSignal) then + FSignal := Null; end; function TSignal.Subscribe(Subscriber: ISubscriber): TSubscription; diff --git a/Src/Myc.TaskManager.pas b/Src/Myc.TaskManager.pas index 83a5c9b..de2a680 100644 --- a/Src/Myc.TaskManager.pas +++ b/Src/Myc.TaskManager.pas @@ -7,27 +7,41 @@ uses Myc.Signals; type - IMycTaskManager = interface - // TOD Dokumentation - function CreateTask(const Gate: TState; const Proc: TProc): TSignal.TSubscription; + TTaskManager = record + type + IMycTaskManager = interface + procedure RunTask(const Gate: TState; const Proc: TProc); + procedure WaitFor(const State: TState); + end; + private + FTaskManager: IMycTaskManager; + + public + constructor Create(const ATaskManager: IMycTaskManager); + class operator Implicit(const A: IMycTaskManager): TTaskManager; overload; + class operator Implicit(const A: TTaskManager): IMycTaskManager; overload; + + // Run a task when the gate is opened. + procedure RunTask(const Gate: TState; const Proc: TProc); inline; // Waits for the operation associated with State to complete. // Must not be called from a worker thread of this factory. // After waiting, or if the state is already set, any first stored exception // from any worker thread of this factory will be re-raised in the calling (main) thread. // Raises ETaskException if called from within a worker thread. - procedure WaitFor(const State: TState.IState); + procedure WaitFor(const State: TState); inline; end; var - TaskManager: IMycTaskManager; + TaskManager: TTaskManager; procedure SetupTaskManagerMock; implementation uses - System.Generics.Collections; + System.Generics.Collections, + Myc.Core.Tasks; type TMycExecMock = class(TInterfacedObject, TSignal.ISubscriber) @@ -38,11 +52,11 @@ type function Notify: Boolean; end; - TMycTaskManagerMock = class(TInterfacedObject, IMycTaskManager) + TMycTaskManagerMock = class(TInterfacedObject, TTaskManager.IMycTaskManager) public // IMycTaskManager - function CreateTask(const Gate: TState; const Proc: TProc): TSignal.TSubscription; - procedure WaitFor(const State: TState.IState); + procedure RunTask(const Gate: TState; const Proc: TProc); + procedure WaitFor(const State: TState); constructor Create; end; @@ -72,19 +86,44 @@ begin inherited Create; end; -function TMycTaskManagerMock.CreateTask(const Gate: TState; const Proc: TProc): TSignal.TSubscription; +procedure TMycTaskManagerMock.RunTask(const Gate: TState; const Proc: TProc); begin - Result := Gate.Signal.Subscribe(TMycExecMock.Create(Proc)); + Gate.Signal.Subscribe(TMycExecMock.Create(Proc)); end; -procedure TMycTaskManagerMock.WaitFor(const State: TState.IState); +procedure TMycTaskManagerMock.WaitFor(const State: TState); begin Assert(State.IsSet); end; procedure SetupTaskManagerMock; begin - TaskManager := TMycTaskManagerMock.Create; + TaskManager.Create(TMycTaskManagerMock.Create); +end; + +constructor TTaskManager.Create(const ATaskManager: IMycTaskManager); +begin + FTaskManager := ATaskManager; +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); +end; + +class operator TTaskManager.Implicit(const A: IMycTaskManager): TTaskManager; +begin + Result.Create(A); +end; + +class operator TTaskManager.Implicit(const A: TTaskManager): IMycTaskManager; +begin + Result := A.FTaskManager; end; end. diff --git a/Src/Myc.Trade.DataStream.pas b/Src/Myc.Trade.DataStream.pas index 72d8bc4..eef52b3 100644 --- a/Src/Myc.Trade.DataStream.pas +++ b/Src/Myc.Trade.DataStream.pas @@ -172,29 +172,6 @@ type function ParseFileName(const FileName: string): TAuraDataFile; override; end; - // Implements a data stream that reads from Aura-specific historical data files. - TAuraFileLoader = class(TInterfacedObject) - type - TDataProc = reference to procedure(const Values: TArray>); - - private - class procedure LoadFile( - DataServer: TAuraDataServer; - currFile: TAuraDataFile; - Terminated: TState; - Done: TLatch; - Proc: TDataProc - ); - - public - class function LoadData( - DataServer: TAuraDataServer; - const Symbol: String; - const Terminate: TSignal; - const Proc: TDataProc - ): TState; - end; - implementation uses @@ -413,12 +390,9 @@ begin var capProc := Proc; var terminated := TFlag.CreateObserver(Terminate).State; - var firstFile := FindFirstFile(Symbol); - - var done := TLatch.CreateLatch(1); - Result := done.State; - - TaskManager.CreateTask(firstFile.Done, procedure begin LoadFile(firstFile.Value, terminated, capProc).Subscribe(done); end); + Result := + FindFirstFile(Symbol) + .Chain(function(const FirstFile: TAuraDataFile): TState begin Result := LoadFile(FirstFile, terminated, capProc); end); end; function TAuraDataServer.LoadDataFile(const DataFile: TAuraDataFile): TFuture>>; @@ -431,19 +405,16 @@ begin if not currFile.IsValid or Terminated.IsSet then exit(TState.Null); - var done := TLatch.CreateLatch(1); - Result := done.State; - var data := FCachedFiles.GetOrAdd(currFile.GetFullFileName); - TaskManager.CreateTask( - data.Done, - procedure - begin - Proc(data.Value, Terminated); - LoadFile(currFile.GetNextFile, Terminated, Proc).Subscribe(done); - end - ); + Result := + data.Chain( + function(const Data: TArray>): TState + begin + Proc(Data, Terminated); + Result := LoadFile(currFile.GetNextFile, Terminated, Proc); + end + ); end; { TAuraFileStream } @@ -792,49 +763,4 @@ begin end; end; -class function TAuraFileLoader.LoadData( - DataServer: TAuraDataServer; - const Symbol: String; - const Terminate: TSignal; - const Proc: TDataProc -): TState; -begin - var capProc := Proc; - var terminated := TFlag.CreateObserver(Terminate).State; - - var firstFile := DataServer.FindFirstFile(Symbol); - - var done := TLatch.CreateLatch(1); - - TaskManager.CreateTask(firstFile.Done, procedure begin LoadFile(DataServer, firstFile.Value, terminated, done, capProc); end); - - Result := done.State; -end; - -class procedure TAuraFileLoader.LoadFile( - DataServer: TAuraDataServer; - currFile: TAuraDataFile; - Terminated: TState; - Done: TLatch; - Proc: TDataProc -); -begin - if not currFile.IsValid or Terminated.IsSet then - begin - Done.Notify; - exit; - end; - - var data := DataServer.LoadDataFile(currFile); - - TaskManager.CreateTask( - data.Done, - procedure - begin - Proc(data.Value); - LoadFile(DataServer, currFile.GetNextFile, Terminated, Done, Proc); - end - ); -end; - end. diff --git a/Test/TestCoreFutures.pas b/Test/TestCoreFutures.pas index b46e7d2..65061cc 100644 --- a/Test/TestCoreFutures.pas +++ b/Test/TestCoreFutures.pas @@ -27,7 +27,7 @@ type [IgnoreMemoryLeaks(true)] TTestMycGateFuncFuture = class(TObject) private - FTaskFactory: IMycTaskFactory; + FTaskFactory: IMycThreadPool; FProcExecutionCount: Integer; // Counter for side effects of AProc FSharedCounter: Integer; // For Fan-Out test side effects public @@ -156,7 +156,7 @@ end; procedure TTestMycGateFuncFuture.Test_ExceptionInProc_HandledAsPlanned; var LFuture: TFuture.IFuture; - LLocalTaskFactory: IMycTaskFactory; + LLocalTaskFactory: IMycThreadPool; LInitStateAsState: TState.IState; LResultValue: Integer; LExpectedExceptionRaisedByFactory: Boolean; diff --git a/Test/TestTasks.pas b/Test/TestTasks.pas index 3b0fa79..36e982e 100644 --- a/Test/TestTasks.pas +++ b/Test/TestTasks.pas @@ -17,7 +17,7 @@ type [IgnoreMemoryLeaks(true)] TMycTaskFactoryTests = class(TObject) private - FFactory: IMycTaskFactory; + FFactory: IMycThreadPool; public [Setup] procedure Setup;