Skip to content

Commit 3c042b8

Browse files
committed
Tradier
1 parent aa32a73 commit 3c042b8

11 files changed

Lines changed: 906 additions & 12 deletions

File tree

Gateways/Schwab/Libs/Schwab.csproj

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,6 @@
88
</PropertyGroup>
99

1010
<ItemGroup>
11-
<PackageReference Include="Flurl.Http" Version="4.0.2" />
1211
<PackageReference Include="SchwabBroker" Version="1.0.2" />
1312
</ItemGroup>
1413

Lines changed: 137 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,137 @@
1+
using Core.Conventions;
2+
using Core.Enums;
3+
using Core.Extensions;
4+
using Core.Grains;
5+
using Core.Models;
6+
using System;
7+
using System.Linq;
8+
using System.Threading;
9+
using System.Threading.Tasks;
10+
using Tradier.Messages.Stream;
11+
using Tradier.Models;
12+
13+
namespace Tradier.Grains
14+
{
15+
public interface ITradierConnectionGrain : IConnectionGrain
16+
{
17+
/// <summary>
18+
/// Connect
19+
/// </summary>
20+
/// <param name="connection"></param>
21+
/// <param name="observer"></param>
22+
Task<StatusResponse> Setup(Connection connection, ITradeObserver observer);
23+
}
24+
25+
/// <summary>
26+
/// Constructor
27+
/// </summary>
28+
public class TradierConnectionGrain : ConnectionGrain, ITradierConnectionGrain
29+
{
30+
/// <summary>
31+
/// State
32+
/// </summary>
33+
protected Connection state;
34+
35+
/// <summary>
36+
/// Connector
37+
/// </summary>
38+
protected TradierBroker connector;
39+
40+
/// <summary>
41+
/// Observer
42+
/// </summary>
43+
protected ITradeObserver observer;
44+
45+
/// <summary>
46+
/// Connect
47+
/// </summary>
48+
/// <param name="connection"></param>
49+
/// <param name="grainObserver"></param>
50+
public virtual async Task<StatusResponse> Setup(Connection connection, ITradeObserver grainObserver)
51+
{
52+
var cleaner = new CancellationTokenSource(connection.Timeout);
53+
54+
await Disconnect();
55+
56+
state = connection;
57+
observer = grainObserver;
58+
connector = new()
59+
{
60+
Token = connection.AccessToken,
61+
SessionToken = connection.SessionToken,
62+
};
63+
64+
await connector.Connect(cleaner.Token);
65+
await Task.WhenAll(connection.Account.Instruments.Values.Select(Subscribe));
66+
67+
return new()
68+
{
69+
Data = StatusEnum.Active
70+
};
71+
}
72+
73+
/// <summary>
74+
/// Save state and dispose
75+
/// </summary>
76+
public override Task<StatusResponse> Disconnect()
77+
{
78+
connections?.ForEach(o => o.Dispose());
79+
connections?.Clear();
80+
connector?.Dispose();
81+
82+
return Task.FromResult(new StatusResponse
83+
{
84+
Data = StatusEnum.Inactive
85+
});
86+
}
87+
88+
/// <summary>
89+
/// Subscribe to streams
90+
/// </summary>
91+
/// <param name="instrument"></param>
92+
public override async Task<StatusResponse> Subscribe(Instrument instrument)
93+
{
94+
var descriptor = this.GetDescriptor();
95+
var instrumentDescriptor = this.GetDescriptor(instrument.Name);
96+
var domGrain = GrainFactory.GetGrain<IDomGrain>(instrumentDescriptor);
97+
var instrumentGrain = GrainFactory.GetGrain<IInstrumentGrain>(instrumentDescriptor);
98+
var positionsGrain = GrainFactory.GetGrain<IPositionsGrain>(descriptor);
99+
var ordersGrain = GrainFactory.GetGrain<IOrdersGrain>(descriptor);
100+
101+
connector.OnPrice += async o =>
102+
{
103+
var group = await instrumentGrain.Send(instrument with
104+
{
105+
Price = MapPrice(o)
106+
});
107+
108+
observer.StreamPrice(group);
109+
110+
await ordersGrain.Tap(group);
111+
await positionsGrain.Tap(group);
112+
await observer.StreamTrade(group);
113+
};
114+
115+
await connector.Subscribe(instrument.Name);
116+
117+
return new()
118+
{
119+
Data = StatusEnum.Active
120+
};
121+
}
122+
123+
/// <summary>
124+
/// Get price
125+
/// </summary>
126+
/// <param name="message"></param>
127+
protected virtual Price MapPrice(PriceMessage message) => new()
128+
{
129+
Ask = message.Ask,
130+
Bid = message.Bid,
131+
AskSize = message.AskSize,
132+
BidSize = message.BidSize,
133+
Last = (message.Bid + message.Ask) / 2.0,
134+
Time = message?.BidDate ?? DateTime.Now.Ticks
135+
};
136+
}
137+
}
Lines changed: 185 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,185 @@
1+
using Core.Enums;
2+
using Core.Grains;
3+
using Core.Models;
4+
using System.Linq;
5+
using System.Threading;
6+
using System.Threading.Tasks;
7+
using Tradier.Messages.MarketData;
8+
using Tradier.Models;
9+
using Tradier.Queries.MarketData;
10+
11+
namespace Tradier.Grains
12+
{
13+
public interface ITradierOptionsGrain : IOptionsGrain
14+
{
15+
/// <summary>
16+
/// Connect
17+
/// </summary>
18+
/// <param name="connection"></param>
19+
Task<StatusResponse> Setup(Connection connection);
20+
}
21+
22+
public class TradierOptionsGrain : OptionsGrain, ITradierOptionsGrain
23+
{
24+
/// <summary>
25+
/// State
26+
/// </summary>
27+
protected Connection state;
28+
29+
/// <summary>
30+
/// Connector
31+
/// </summary>
32+
protected TradierBroker connector;
33+
34+
/// <summary>
35+
/// Connect
36+
/// </summary>
37+
/// <param name="connection"></param>
38+
public virtual async Task<StatusResponse> Setup(Connection connection)
39+
{
40+
var cleaner = new CancellationTokenSource(connection.Timeout);
41+
42+
state = connection;
43+
connector = new()
44+
{
45+
Token = connection.AccessToken,
46+
SessionToken = connection.SessionToken,
47+
};
48+
49+
await connector.Connect(cleaner.Token);
50+
51+
return new()
52+
{
53+
Data = StatusEnum.Active
54+
};
55+
}
56+
57+
/// <summary>
58+
/// Option chain
59+
/// </summary>
60+
/// <param name="criteria"></param>
61+
public override async Task<InstrumentsResponse> Options(Criteria criteria)
62+
{
63+
var query = new OptionChainRequest
64+
{
65+
Symbol = criteria.Instrument.Name,
66+
Expiration = criteria.MaxDate ?? criteria.MinDate
67+
};
68+
69+
var cleaner = new CancellationTokenSource(state.Timeout);
70+
var chain = await connector.GetOptionChain(query, cleaner.Token);
71+
var options = chain
72+
.Options
73+
?.Select(MapOption)
74+
?.OrderBy(o => o.Derivative.ExpirationDate)
75+
?.ThenBy(o => o.Derivative.Strike)
76+
?.ThenBy(o => o.Derivative.Side)
77+
?.ToList() ?? [];
78+
79+
return new()
80+
{
81+
Data = options
82+
};
83+
}
84+
85+
/// <summary>
86+
/// Get internal option
87+
/// </summary>
88+
/// <param name="message"></param>
89+
protected virtual Instrument MapOption(OptionMessage message)
90+
{
91+
var instrument = new Instrument
92+
{
93+
Name = message.Underlying,
94+
Exchange = message.Exchange
95+
};
96+
97+
var optionPoint = new Price
98+
{
99+
Ask = message.Ask,
100+
Bid = message.Bid,
101+
AskSize = message.AskSize ?? 0,
102+
BidSize = message.BidSize ?? 0,
103+
Volume = message.Volume,
104+
Last = message.Last,
105+
Bar = new()
106+
{
107+
Low = message.Low,
108+
High = message.High,
109+
Open = message.Open,
110+
Close = message.Close
111+
}
112+
};
113+
114+
var optionInstrument = new Instrument
115+
{
116+
Basis = instrument,
117+
Price = optionPoint,
118+
Name = message.Symbol,
119+
Exchange = message.Exchange,
120+
Leverage = message.ContractSize ?? 100,
121+
Type = GetInstrumentType(message.Type)
122+
};
123+
124+
var side = null as OptionSideEnum?;
125+
126+
switch (message.OptionType.ToUpper())
127+
{
128+
case "PUT": side = OptionSideEnum.Put; break;
129+
case "CALL": side = OptionSideEnum.Call; break;
130+
}
131+
132+
var derivative = new Derivative
133+
{
134+
Side = side,
135+
Strike = message.Strike,
136+
TradeDate = message.ExpirationDate,
137+
ExpirationDate = message.ExpirationDate,
138+
OpenInterest = message.OpenInterest ?? 0,
139+
Volatility = message?.Greeks?.SmvIV ?? 0,
140+
};
141+
142+
var greeks = message?.Greeks;
143+
144+
if (greeks is not null)
145+
{
146+
derivative = derivative with
147+
{
148+
Variance = new Variance
149+
{
150+
Rho = greeks.Rho ?? 0,
151+
Vega = greeks.Vega ?? 0,
152+
Delta = greeks.Delta ?? 0,
153+
Gamma = greeks.Gamma ?? 0,
154+
Theta = greeks.Theta ?? 0
155+
}
156+
};
157+
}
158+
159+
optionInstrument = optionInstrument with
160+
{
161+
Price = optionPoint,
162+
Derivative = derivative
163+
};
164+
165+
return optionInstrument;
166+
}
167+
168+
/// <summary>
169+
/// Asset type
170+
/// </summary>
171+
/// <param name="assetType"></param>
172+
protected virtual InstrumentEnum? GetInstrumentType(string assetType)
173+
{
174+
switch (assetType?.ToUpper())
175+
{
176+
case "EQUITY": return InstrumentEnum.Shares;
177+
case "INDEX": return InstrumentEnum.Indices;
178+
case "FUTURE": return InstrumentEnum.Futures;
179+
case "OPTION": return InstrumentEnum.Options;
180+
}
181+
182+
return InstrumentEnum.Group;
183+
}
184+
}
185+
}
Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
1+
using Core.Enums;
2+
using Core.Models;
3+
using System.Linq;
4+
using System.Threading;
5+
using System.Threading.Tasks;
6+
7+
namespace Tradier.Grains
8+
{
9+
public interface ITradierOrderSenderGrain : ITradierOrdersGrain
10+
{
11+
}
12+
13+
public class TradierOrderSenderGrain : TradierOrdersGrain, ITradierOrderSenderGrain
14+
{
15+
/// <summary>
16+
/// Send order
17+
/// </summary>
18+
/// <param name="order"></param>
19+
public override async Task<OrderResponse> Send(Order order)
20+
{
21+
//var message = MapOrder(order);
22+
//var accountCode = order.Account.Descriptor;
23+
//var messageResponse = await connector.SendOrder(message, accountCode, CancellationToken.None);
24+
//var response = new OrderResponse { Data = new() { Id = messageResponse.OrderId } };
25+
26+
return null;
27+
}
28+
}
29+
}

0 commit comments

Comments
 (0)