///|
pub(all) enum MarketSourceAvailability {
LocalFixture
ExternalPlanned
} derive(Debug, Eq, ToJson, FromJson)
///|
pub(all) enum MarketFetchStatus {
SnapshotReady
AwaitingExternalSnapshot
SnapshotRejected
} derive(Debug, Eq, ToJson, FromJson)
///|
pub(all) enum MarketDataIssueKind {
SourceMismatch
SymbolMismatch
TimeframeMismatch
EmptySnapshot
OpenBarPresent
InsufficientClosedBars
} derive(Debug, Eq, ToJson, FromJson)
///|
pub(all) struct MarketSourceRequest {
data_source : @domain.DataSource
symbol : @domain.Symbol
timeframe : @domain.Timeframe
limit : Int
captured_at_ms : Int64
} derive(Debug, Eq, ToJson, FromJson)
///|
pub(all) struct MarketSourcePlan {
request : MarketSourceRequest
source_id : @domain.SourceId
adapter : String
raw_path : String
availability : MarketSourceAvailability
} derive(Debug, Eq, ToJson, FromJson)
///|
pub(all) struct MarketFetchResult {
plan : MarketSourcePlan
snapshot : @domain.KlineSnapshot?
status : MarketFetchStatus
issues : Array[MarketDataIssue]
summary : String
} derive(Debug, Eq, ToJson, FromJson)
///|
pub(all) struct MarketDataIssue {
kind : MarketDataIssueKind
path : String
message : String
} derive(Debug, Eq, ToJson, FromJson)
///|
pub(all) struct MarketIngestInput {
request : MarketSourceRequest
snapshot : @domain.KlineSnapshot
} derive(Debug, Eq, ToJson, FromJson)
///|
fn normalized_limit(limit : Int) -> Int {
if limit < 1 {
1
} else {
limit
}
}
///|
fn data_source_id(source : @domain.DataSource) -> String {
match source {
TradingView => "tradingview"
Mt5 => "mt5"
Yfinance => "yfinance"
AkShare => "akshare"
Fixture => "fixture"
}
}
///|
fn adapter_name(source : @domain.DataSource) -> String {
match source {
TradingView => "planned-tradingview-adapter"
Mt5 => "planned-mt5-adapter"
Yfinance => "planned-yfinance-adapter"
AkShare => "planned-akshare-adapter"
Fixture => "local-fixture-adapter"
}
}
///|
pub fn market_source_request(
data_source : @domain.DataSource,
symbol : @domain.Symbol,
timeframe : @domain.Timeframe,
limit? : Int = 60,
captured_at_ms? : Int64 = 0L,
) -> MarketSourceRequest {
{
data_source,
symbol,
timeframe,
limit: normalized_limit(limit),
captured_at_ms,
}
}
///|
pub fn market_data_issue(
kind : MarketDataIssueKind,
path : String,
message : String,
) -> MarketDataIssue {
{ kind, path, message }
}
///|
pub fn market_ingest_input(
request : MarketSourceRequest,
snapshot : @domain.KlineSnapshot,
) -> MarketIngestInput {
{ request, snapshot }
}
///|
pub fn plan_market_source(request : MarketSourceRequest) -> MarketSourcePlan {
let source_label = data_source_id(request.data_source)
let source_id = "\{source_label}:\{request.symbol}:\{request.timeframe}:\{request.limit}"
{
request,
source_id,
adapter: adapter_name(request.data_source),
raw_path: "raw/market-data/\{source_id}.json",
availability: if request.data_source is Fixture {
LocalFixture
} else {
ExternalPlanned
},
}
}
///|
fn fetch_summary(status : MarketFetchStatus, plan : MarketSourcePlan) -> String {
match status {
SnapshotReady => "market snapshot ready from \{plan.adapter}"
AwaitingExternalSnapshot =>
"awaiting external market snapshot from \{plan.adapter}"
SnapshotRejected => "market snapshot rejected for \{plan.adapter}"
}
}
///|
fn market_fetch_result(
plan : MarketSourcePlan,
snapshot : @domain.KlineSnapshot?,
status : MarketFetchStatus,
issues? : Array[MarketDataIssue] = [],
) -> MarketFetchResult {
{ plan, snapshot, status, issues, summary: fetch_summary(status, plan) }
}
///|
fn fixture_bar(index : Int, limit : Int) -> @domain.KlineBar {
let seq = limit - index
let base = 100.0 + index.to_double()
@domain.kline_bar(
seq,
base,
base + 2.0,
base - 1.0,
base + 1.0,
volume=1000.0 + index.to_double(),
closed=true,
)
}
///|
fn fixture_bars(limit : Int) -> Array[@domain.KlineBar] {
let bars : Array[@domain.KlineBar] = []
for index in 0.. MarketFetchResult {
let plan = plan_market_source(request)
if request.data_source is Fixture {
market_fetch_result(
plan,
Some(
@domain.kline_snapshot(
plan.source_id,
request.data_source,
request.symbol,
request.timeframe,
request.captured_at_ms,
fixture_bars(request.limit),
),
),
SnapshotReady,
)
} else {
market_fetch_result(plan, None, AwaitingExternalSnapshot)
}
}
///|
fn collect_ingest_issues(
request : MarketSourceRequest,
snapshot : @domain.KlineSnapshot,
) -> Array[MarketDataIssue] {
let issues : Array[MarketDataIssue] = []
if snapshot.data_source != request.data_source {
issues.push(
market_data_issue(
SourceMismatch,
"snapshot.data_source",
"snapshot source does not match requested source",
),
)
}
if snapshot.symbol != request.symbol {
issues.push(
market_data_issue(
SymbolMismatch,
"snapshot.symbol",
"snapshot symbol does not match requested symbol",
),
)
}
if snapshot.timeframe != request.timeframe {
issues.push(
market_data_issue(
TimeframeMismatch,
"snapshot.timeframe",
"snapshot timeframe does not match requested timeframe",
),
)
}
if snapshot.bars.is_empty() {
issues.push(
market_data_issue(EmptySnapshot, "snapshot.bars", "snapshot has no bars"),
)
}
if snapshot.bars.any(fn(bar) { !bar.closed }) {
issues.push(
market_data_issue(
OpenBarPresent,
"snapshot.bars.closed",
"snapshot includes open bars; Moonfish requires closed-bar evidence",
),
)
}
if snapshot.closed_bar_count() < request.limit {
issues.push(
market_data_issue(
InsufficientClosedBars,
"snapshot.bars",
"snapshot has fewer closed bars than requested limit",
),
)
}
issues
}
///|
pub fn ingest_market_snapshot(input : MarketIngestInput) -> MarketFetchResult {
let plan = plan_market_source(input.request)
let issues = collect_ingest_issues(input.request, input.snapshot)
if issues.is_empty() {
market_fetch_result(plan, Some(input.snapshot), SnapshotReady)
} else {
market_fetch_result(plan, None, SnapshotRejected, issues~)
}
}