diff --git a/AuraTrader/MainForm.pas b/AuraTrader/MainForm.pas index c1ef0c0..6a00f60 100644 --- a/AuraTrader/MainForm.pas +++ b/AuraTrader/MainForm.pas @@ -89,15 +89,14 @@ begin exit; var arr := TDataSeries.CreateArray(3000); - var curr := TMutable>.CreateProtected; + var curr := TWriteable>.CreateProtected; var Symbol := FSymbols.WaitFor[SymbolsComboBox.ItemIndex]; var done := - TAuraFileLoader.LoadData( - FServer as TAuraTABFileServer, + FServer.ProcessData( Symbol, FTerminate.Signal, - procedure(const Values: TArray>) + procedure(const Values: TArray>; const Terminated: TState) begin arr.Add(Values); curr.Value := arr.Copy; diff --git a/Src/Myc.Core.Lazy.pas b/Src/Myc.Core.Lazy.pas index 676d3d2..e93d468 100644 --- a/Src/Myc.Core.Lazy.pas +++ b/Src/Myc.Core.Lazy.pas @@ -36,7 +36,7 @@ type constructor Create(const AChanged: TSignal; const AProc: TFunc); end; - TMycWriteableMutable = class(TInterfacedObject, TMutable.IMutable, TMutable.IWriteable) + TMycWriteableMutable = class(TInterfacedObject, TMutable.IMutable, TWriteable.IWriteable) private FValue: T; FChanged: TEvent; @@ -77,7 +77,7 @@ type constructor Create(const AChanged: TSignal.ISignal; const AProc: TFunc); end; - TMycProtectedWriteableMutable = class(TInterfacedObject, TMutable.IMutable, TMutable.IWriteable) + TMycProtectedWriteableMutable = class(TInterfacedObject, TMutable.IMutable, TWriteable.IWriteable) private FValue: T; FChanged: TEvent; @@ -85,6 +85,8 @@ type protected function GetChanged: TSignal; function GetValue: T; + procedure Lock; inline; + procedure Release; inline; public constructor Create(const AValue: T); procedure SetValue(const Value: T); @@ -234,24 +236,33 @@ end; function TMycProtectedWriteableMutable.GetValue: T; begin - while AtomicExchange(FLock, 1) = 1 do - YieldProcessor; + Lock; try Result := FValue; finally - AtomicExchange(FLock, 0); + Release; end; end; +procedure TMycProtectedWriteableMutable.Lock; +begin + while AtomicExchange(FLock, 1) = 1 do + YieldProcessor; +end; + +procedure TMycProtectedWriteableMutable.Release; +begin + AtomicExchange(FLock, 0); +end; + procedure TMycProtectedWriteableMutable.SetValue(const Value: T); begin - while AtomicExchange(FLock, 1) = 1 do - YieldProcessor; + Lock; try FValue := Value; FChanged.Notify; finally - AtomicExchange(FLock, 0); + Release; end; end; diff --git a/Src/Myc.Lazy.pas b/Src/Myc.Lazy.pas index 5fffd26..55a2cfe 100644 --- a/Src/Myc.Lazy.pas +++ b/Src/Myc.Lazy.pas @@ -18,11 +18,6 @@ type property Value: T read GetValue; end; - IWriteable = interface(IMutable) - procedure SetValue(const Value: T); - property Value: T read GetValue write SetValue; - end; - {$REGION 'private'} strict private class var @@ -44,13 +39,36 @@ type class function Construct(const Changing: TSignal; const Proc: TFunc): TMutable; overload; static; - class function CreateWriteable(const Init: T): IWriteable; overload; static; - class function CreateProtected: IWriteable; overload; static; - property Value: T read GetValue; property Changed: TSignal read GetChanged; end; + TWriteable = record + type + IWriteable = interface(TMutable.IMutable) + procedure SetValue(const Value: T); + property Value: T read GetValue write SetValue; + end; + + private + FWriteable: IWriteable; + function GetChanged: TSignal; inline; + function GetValue: T; inline; + procedure SetValue(const Value: T); inline; + + public + constructor Create(const AWriteable: IWriteable); + + class function CreateProtected: TWriteable; overload; static; + class function CreateWriteable: TWriteable; overload; static; + + class operator Implicit(const A: IWriteable): TWriteable; overload; + class operator Implicit(const A: TWriteable): IWriteable; overload; + + property Value: T read GetValue write SetValue; + property Changed: TSignal read GetChanged; + end; + TLazy = record type ILazy = interface @@ -109,16 +127,6 @@ begin Result := TMycFuncMutable.Create(Changing, Proc); end; -class function TMutable.CreateProtected: IWriteable; -begin - Result := TMycProtectedWriteableMutable.Create(Default(T)); -end; - -class function TMutable.CreateWriteable(const Init: T): IWriteable; -begin - Result := TMycWriteableMutable.Create(Init); -end; - function TMutable.GetChanged: TSignal; begin Result := FMutable.Changed; @@ -193,4 +201,46 @@ begin Dest.FLazy := FNull; end; +{ TWriteable } + +constructor TWriteable.Create(const AWriteable: IWriteable); +begin + FWriteable := AWriteable; +end; + +class function TWriteable.CreateProtected: TWriteable; +begin + Result := TMycProtectedWriteableMutable.Create(Default(T)); +end; + +class function TWriteable.CreateWriteable: TWriteable; +begin + Result := TMycWriteableMutable.Create(Default(T)); +end; + +function TWriteable.GetChanged: TSignal; +begin + Result := FWriteable.Changed; +end; + +function TWriteable.GetValue: T; +begin + Result := FWriteable.Value; +end; + +procedure TWriteable.SetValue(const Value: T); +begin + FWriteable.Value := Value; +end; + +class operator TWriteable.Implicit(const A: IWriteable): TWriteable; +begin + Result.Create(A); +end; + +class operator TWriteable.Implicit(const A: TWriteable): IWriteable; +begin + Result := A.FWriteable; +end; + end. diff --git a/Src/Myc.Trade.DataStream.pas b/Src/Myc.Trade.DataStream.pas index f34afb5..72d8bc4 100644 --- a/Src/Myc.Trade.DataStream.pas +++ b/Src/Myc.Trade.DataStream.pas @@ -54,23 +54,14 @@ type function IsHistory: Boolean; virtual; abstract; end; - // Represents a factory for creating IDataStream instances. - IDataStreamNode = interface - ['{DA3531F1-158E-4A63-8E5A-34089C537B1B}'] - function CreateDataStream: IDataStream; - end; - - // Abstract base class for IDataStreamNode implementations. - TDataStreamNode = class(TInterfacedObject, IDataStreamNode) - public - function CreateDataStream: IDataStream; virtual; abstract; - end; + TDataProc = reference to procedure(const Values: TArray>; const Terminated: TState); // Interface for an instantiable data server. IDataServer = interface ['{1F8E5A9D-E92A-44C1-9F3F-C4B82A6E94B3}'] function CreateStream(const Symbol: String): IDataStream; procedure ClearCache; + function ProcessData(const Symbol: String; const Terminate: TSignal; const Proc: TDataProc): TState; function EnumerateSymbols: TFuture>; end; @@ -110,6 +101,8 @@ type // Used by cache to actually load an uncached file. function DoLoad(const FileName: string): TFuture>>; + function LoadFile(currFile: TAuraDataFile; Terminated: TState; Proc: TDataProc): TState; + strict private class var FLoadGate: TLatch; @@ -130,6 +123,9 @@ type function ParseFileName(const FileName: string): TAuraDataFile; virtual; abstract; function EnumerateSymbols: TFuture>; function LoadDataFile(const DataFile: TAuraDataFile): TFuture>>; + + function ProcessData(const Symbol: String; const Terminate: TSignal; const Proc: TDataProc): TState; + property Path: String read GetPath; end; @@ -180,13 +176,8 @@ type TAuraFileLoader = class(TInterfacedObject) type TDataProc = reference to procedure(const Values: TArray>); - public - class function LoadData( - DataServer: TAuraDataServer; - const Symbol: String; - const Terminate: TSignal; - const Proc: TDataProc - ): TState; + + private class procedure LoadFile( DataServer: TAuraDataServer; currFile: TAuraDataFile; @@ -194,6 +185,14 @@ type Done: TLatch; Proc: TDataProc ); + + public + class function LoadData( + DataServer: TAuraDataServer; + const Symbol: String; + const Terminate: TSignal; + const Proc: TDataProc + ): TState; end; implementation @@ -409,11 +408,44 @@ begin Result := FPath; end; +function TAuraDataServer.ProcessData(const Symbol: String; const Terminate: TSignal; const Proc: TDataProc): TState; +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); +end; + function TAuraDataServer.LoadDataFile(const DataFile: TAuraDataFile): TFuture>>; begin Result := FCachedFiles.GetOrAdd(DataFile.GetFullFileName); end; +function TAuraDataServer.LoadFile(currFile: TAuraDataFile; Terminated: TState; Proc: TDataProc): TState; +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 + ); +end; + { TAuraFileStream } constructor TAuraFileStream.Create(ADataServer: TAuraDataServer; const ASymbol: String);