use std::collections::BTreeSet;
use std::path::Path;
use pomelo_data::{
list_symbols, load_panel, write_combined_panel, Field, LocalSource, ObjectSink, PANELS_DIR,
PRICES_DIR,
};
use serde_json::Value;
use yuzu_core::panel::Panel;
use super::config::SyncConfig;
use super::http::Fetcher;
use super::util::iso_to_i32;
use super::HttpClient;
use super::FMP_BASE;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Index {
Sp500,
Nasdaq,
DowJones,
}
pub const MEMBERSHIP_SERIES: &[&str] = &["in_sp500", "in_nasdaq", "in_dowjones"];
impl Index {
pub fn parse(s: &str) -> Option<Index> {
match s.trim().to_ascii_lowercase().as_str() {
"sp500" | "sp-500" | "spx" | "spy" => Some(Index::Sp500),
"nasdaq" | "ndx" | "nasdaq100" => Some(Index::Nasdaq),
"dowjones" | "dow" | "djia" | "dji" => Some(Index::DowJones),
_ => None,
}
}
fn current_endpoint(self) -> &'static str {
match self {
Index::Sp500 => "sp-500",
Index::Nasdaq => "nasdaq",
Index::DowJones => "dow-jones",
}
}
fn historical_endpoint(self) -> &'static str {
match self {
Index::Sp500 => "historical-sp-500",
Index::Nasdaq => "historical-nasdaq",
Index::DowJones => "historical-dow-jones",
}
}
pub fn series_name(self) -> &'static str {
match self {
Index::Sp500 => "in_sp500",
Index::Nasdaq => "in_nasdaq",
Index::DowJones => "in_dowjones",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Change {
pub date: i32,
pub added: Option<String>,
pub removed: Option<String>,
}
fn nonempty(obj: &serde_json::Map<String, Value>, key: &str) -> Option<String> {
obj.get(key)
.and_then(Value::as_str)
.map(str::trim)
.filter(|s| !s.is_empty())
.map(str::to_string)
}
pub(crate) fn parse_current(rows: &[Value]) -> BTreeSet<String> {
rows.iter()
.filter_map(|r| nonempty(r.as_object()?, "symbol"))
.collect()
}
pub(crate) fn parse_changes(rows: &[Value]) -> Vec<Change> {
let mut out: Vec<Change> = rows
.iter()
.filter_map(|r| {
let obj = r.as_object()?;
let date = obj
.get("date")
.and_then(Value::as_str)
.and_then(iso_to_i32)?;
let added = nonempty(obj, "symbol");
let removed = nonempty(obj, "removedTicker");
if added.is_none() && removed.is_none() {
return None;
}
Some(Change {
date,
added,
removed,
})
})
.collect();
out.sort_by_key(|c| c.date);
out
}
pub struct IndexMembership {
index: Index,
current: BTreeSet<String>,
changes: Vec<Change>,
}
impl IndexMembership {
pub fn fetch<H: HttpClient>(
http: &H,
api_key: &str,
index: Index,
cfg: &SyncConfig,
) -> Result<IndexMembership, String> {
let fetcher = Fetcher::new(http, cfg);
let cur = fetcher.get_rows(&Self::url(index.current_endpoint(), api_key))?;
let hist = fetcher.get_rows(&Self::url(index.historical_endpoint(), api_key))?;
let current = parse_current(&cur);
if current.is_empty() {
return Err(format!(
"index '{}' returned no current constituents",
index.series_name()
));
}
Ok(IndexMembership {
index,
current,
changes: parse_changes(&hist),
})
}
fn url(endpoint: &str, key: &str) -> String {
format!("{FMP_BASE}/stable/{endpoint}?apikey={key}")
}
pub fn series_name(&self) -> &'static str {
self.index.series_name()
}
pub(crate) fn members_asof(&self, date: i32) -> BTreeSet<String> {
let mut set = self.current.clone();
for c in self.changes.iter().rev() {
if c.date <= date {
break;
}
if let Some(a) = &c.added {
set.remove(a);
}
if let Some(r) = &c.removed {
set.insert(r.clone());
}
}
set
}
pub fn ever_members(&self, from: i32, to: i32) -> Vec<String> {
let mut set = self.members_asof(from);
for c in &self.changes {
if c.date < from || c.date > to {
continue;
}
if let Some(a) = &c.added {
set.insert(a.clone());
}
if let Some(r) = &c.removed {
set.insert(r.clone());
}
}
set.into_iter().collect()
}
pub fn membership_panel(&self, calendar: &[i32], columns: &[String]) -> Result<Panel, String> {
if calendar.is_empty() {
return Err("empty calendar for membership panel".to_string());
}
let mut set = self.members_asof(calendar[0]);
let mut p = self.changes.partition_point(|c| c.date <= calendar[0]);
let mut rows: Vec<Vec<f64>> = Vec::with_capacity(calendar.len());
for &day in calendar {
while p < self.changes.len() && self.changes[p].date <= day {
let c = &self.changes[p];
if let Some(a) = &c.added {
set.insert(a.clone());
}
if let Some(r) = &c.removed {
set.remove(r);
}
p += 1;
}
rows.push(
columns
.iter()
.map(|s| if set.contains(s) { 1.0 } else { 0.0 })
.collect(),
);
}
Panel::from_rows(calendar.to_vec(), columns.to_vec(), rows).map_err(|e| e.to_string())
}
}
fn trading_calendar(root: &Path, from: i32, to: i32) -> Result<Vec<i32>, String> {
let syms = list_symbols(root).map_err(|e| e.to_string())?;
if syms.is_empty() {
return Err("no synced prices to derive a trading calendar from".to_string());
}
let close = load_panel(
&LocalSource::new(root),
&syms,
Field::AdjClose,
from,
to,
PRICES_DIR,
)
.map_err(|e| e.to_string())?;
Ok(close.dates)
}
pub fn write_index_membership(
root: &Path,
membership: &IndexMembership,
from: i32,
to: i32,
) -> Result<(usize, usize), String> {
let calendar = trading_calendar(root, from, to)?;
let columns = membership.ever_members(from, to);
let panel = membership.membership_panel(&calendar, &columns)?;
let bytes = write_combined_panel(&panel).map_err(|e| e.to_string())?;
let key = format!("{PANELS_DIR}/{}.csv.gz", membership.series_name());
LocalSource::new(root)
.put(&key, &bytes)
.map_err(|e| e.to_string())?;
Ok((calendar.len(), columns.len()))
}
#[cfg(test)]
mod tests {
use super::*;
fn fixture() -> IndexMembership {
IndexMembership {
index: Index::Sp500,
current: ["AAA", "BBB", "CCC"]
.iter()
.map(|s| s.to_string())
.collect(),
changes: vec![
Change {
date: 20150601,
added: Some("DDD".into()),
removed: Some("AAA".into()),
},
Change {
date: 20200601,
added: Some("AAA".into()),
removed: Some("DDD".into()),
},
],
}
}
fn set(items: &[&str]) -> BTreeSet<String> {
items.iter().map(|s| s.to_string()).collect()
}
#[test]
fn members_asof_replays_the_log_backwards() {
let m = fixture();
assert_eq!(m.members_asof(20990101), set(&["AAA", "BBB", "CCC"]));
assert_eq!(m.members_asof(20180101), set(&["BBB", "CCC", "DDD"]));
assert_eq!(m.members_asof(20100101), set(&["AAA", "BBB", "CCC"]));
assert_eq!(m.members_asof(20200601), set(&["AAA", "BBB", "CCC"]));
assert_eq!(m.members_asof(20200531), set(&["BBB", "CCC", "DDD"]));
}
#[test]
fn ever_members_is_the_windowed_union() {
let m = fixture();
assert_eq!(
m.ever_members(20100101, 20250101),
vec!["AAA", "BBB", "CCC", "DDD"]
);
assert_eq!(
m.ever_members(20210101, 20220101),
vec!["AAA", "BBB", "CCC"]
);
}
#[test]
fn membership_panel_is_per_day_zero_one() {
let m = fixture();
let cols: Vec<String> = ["AAA", "BBB", "CCC", "DDD"]
.iter()
.map(|s| s.to_string())
.collect();
let cal = [20180101, 20200601, 20210101];
let p = m.membership_panel(&cal, &cols).unwrap();
assert_eq!(p.data.row(0).to_vec(), vec![0.0, 1.0, 1.0, 1.0]);
assert_eq!(p.data.row(1).to_vec(), vec![1.0, 1.0, 1.0, 0.0]);
assert_eq!(p.data.row(2).to_vec(), vec![1.0, 1.0, 1.0, 0.0]);
assert_eq!(p.dates, cal.to_vec());
}
#[test]
fn parse_changes_drops_undated_and_sorts_ascending() {
let rows = vec![
serde_json::json!({"date":"2020-06-01","symbol":"AAA","removedTicker":"DDD"}),
serde_json::json!({"date":"2015-06-01","symbol":"DDD","removedTicker":"AAA"}),
serde_json::json!({"symbol":"NODATE","removedTicker":"X"}), serde_json::json!({"date":"2019-01-01","symbol":"","removedTicker":""}), ];
let changes = parse_changes(&rows);
assert_eq!(changes.len(), 2);
assert_eq!(changes[0].date, 20150601); assert_eq!(changes[1].date, 20200601);
assert_eq!(changes[0].added.as_deref(), Some("DDD"));
assert_eq!(changes[0].removed.as_deref(), Some("AAA"));
}
#[test]
fn index_parse_and_series_names_are_consistent() {
assert_eq!(Index::parse("SP500"), Some(Index::Sp500));
assert_eq!(Index::parse("nasdaq"), Some(Index::Nasdaq));
assert_eq!(Index::parse("dow"), Some(Index::DowJones));
assert_eq!(Index::parse("russell2000"), None);
for idx in [Index::Sp500, Index::Nasdaq, Index::DowJones] {
assert!(MEMBERSHIP_SERIES.contains(&idx.series_name()));
}
}
}