mod connection;
mod error;
mod reporter;
pub use crate::connection::{
HEARTBEAT_IVL_MS, HEARTBEAT_TIMEOUT_MS, POLL_INTERVAL_MS, RECONNECT_MAX_MS,
RECONNECT_MIN_MS,
};
pub use crate::error::Error;
use crate::connection::{Connection, Stall};
use crate::reporter::Reporter;
use chrono::prelude::*;
use elite_journal::entry::market::{BlackMarket, Outfitting, Shipyard};
use elite_journal::entry::{Entry, Event, Market};
use miniz_oxide::inflate;
use serde::Deserialize;
use std::thread;
use std::time::Duration;
use tracing::{debug, info, warn, Level};
pub const URL: &'static str = "tcp://eddn.edcd.io:9500";
#[derive(Debug)]
pub struct Envelope {
pub schema_ref: String,
pub header: Header,
pub message: Message,
pub live: bool,
pub version: Option<String>,
}
#[derive(Deserialize)]
struct RawEnvelope {
#[serde(rename = "$schemaRef")]
schema_ref: String,
header: Header,
message: serde_json::Value,
}
impl<'de> Deserialize<'de> for Envelope {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let raw = RawEnvelope::deserialize(deserializer)?;
let (message, live, version) = {
let schema = Schema::read(&raw.schema_ref);
(
Message::read(schema.map(|schema| schema.name), raw.message)
.map_err(serde::de::Error::custom)?,
schema.map(|schema| schema.live).unwrap_or(true),
schema.map(|schema| schema.version.to_owned()),
)
};
Ok(Envelope {
schema_ref: raw.schema_ref,
header: raw.header,
message,
live,
version,
})
}
}
#[derive(Debug, Deserialize)]
pub struct Header {
#[serde(rename = "gatewayTimestamp")]
pub gateway_timestamp: DateTime<Utc>,
#[serde(rename = "softwareName")]
pub software_name: String,
#[serde(rename = "softwareVersion")]
pub software_version: String,
#[serde(rename = "uploaderID")]
pub uploader_id: String,
}
#[derive(Debug)]
pub enum Message {
Journal(Entry<Event>),
Commodity(Entry<Market>),
Outfitting(Entry<Outfitting>),
Shipyard(Entry<Shipyard>),
BlackMarket(Entry<BlackMarket>),
Unmodeled(serde_json::Value),
}
const SCHEMAS: &str = "/schemas/";
const JOURNAL_SCHEMAS: &[&str] = &[
"journal",
"approachsettlement",
"codexentry",
"dockingdenied",
"dockinggranted",
"fssallbodiesfound",
"fssbodysignals",
"fssdiscoveryscan",
"fsssignaldiscovered",
"navbeaconscan",
"navroute",
"scanbarycentre",
];
impl Message {
fn read(
name: Option<&str>,
message: serde_json::Value,
) -> Result<Self, serde_json::Error> {
let Some(name) = name else {
return Ok(Message::Unmodeled(message));
};
Ok(match name {
name if JOURNAL_SCHEMAS.contains(&name) => {
Message::Journal(serde_json::from_value(message)?)
}
"commodity" => Message::Commodity(serde_json::from_value(message)?),
"outfitting" => {
Message::Outfitting(serde_json::from_value(message)?)
}
"shipyard" => Message::Shipyard(serde_json::from_value(message)?),
"blackmarket" => {
Message::BlackMarket(serde_json::from_value(message)?)
}
_ => Message::Unmodeled(message),
})
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct Schema<'a> {
name: &'a str,
version: &'a str,
live: bool,
}
impl<'a> Schema<'a> {
fn read(schema_ref: &'a str) -> Option<Self> {
let (_, tail) = schema_ref.split_once(SCHEMAS)?;
let mut parts = tail.split('/');
let name = parts.next()?;
let version = parts.next()?;
Some(Schema { name, version, live: parts.next() != Some("test") })
}
}
pub fn subscribe(
url: &str,
stall_timeout: Option<Duration>,
) -> EnvelopeIterator {
let ctx = zmq::Context::new();
let connection =
Connection::open(&ctx, url).expect("failed to open socket");
info!("Subscribed to {}", url);
EnvelopeIterator {
ctx,
url: url.to_string(),
connection,
reports: Reporter::default(),
stall: stall_timeout.map(Stall::new),
}
}
pub struct EnvelopeIterator {
ctx: zmq::Context,
url: String,
connection: Connection,
reports: Reporter,
stall: Option<Stall>,
}
impl EnvelopeIterator {
fn reconnect(&mut self, reason: &str) {
warn!("{}, replacing the connection", reason);
loop {
match Connection::open(&self.ctx, &self.url) {
Ok(connection) => {
self.connection = connection;
if let Some(stall) = &mut self.stall {
stall.restart();
}
self.reports.replaced();
return;
}
Err(err) => {
warn!("Could not open a socket: {}", err);
thread::sleep(Duration::from_secs(5));
}
}
}
}
}
fn read_frame(compressed: &[u8]) -> Option<Result<Envelope, Error>> {
let read = inflate::decompress_to_vec_zlib(compressed)
.map_err(Error::Decompress)
.and_then(|json| {
serde_json::from_slice::<Envelope>(&json).map_err(|source| {
Error::Parse {
source,
json: String::from_utf8_lossy(&json).into_owned(),
}
})
});
if let Ok(envelope) = &read {
if !envelope.live {
debug!(schema = %envelope.schema_ref, "not live");
return None;
}
}
Some(read)
}
impl Iterator for EnvelopeIterator {
type Item = Result<Envelope, Error>;
fn next(&mut self) -> Option<Self::Item> {
loop {
let events = self.connection.events();
for note in self.reports.observe(&events) {
if note.level == Level::WARN {
warn!("{}", note.message);
} else {
info!("{}", note.message);
}
}
match self.connection.socket.recv_bytes(0) {
Ok(compressed) => {
if let Some(stall) = &mut self.stall {
stall.restart();
}
if let Some(read) = read_frame(&compressed) {
return Some(read);
}
}
Err(zmq::Error::EAGAIN) => {
let overrun = self.stall.as_ref().and_then(Stall::overrun);
if let Some(quiet) = overrun {
let reason =
format!("nothing for {}s", quiet.as_secs());
self.reconnect(&reason);
}
}
Err(zmq::Error::EINTR) => {}
Err(err) => {
self.reconnect(&format!("socket error: {}", err));
return Some(Err(Error::Socket(err)));
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
pub(super) fn envelope_json(schema_ref: &str, message: &str) -> String {
format!(
r#"{{
"$schemaRef": "{}",
"header": {{
"gatewayTimestamp": "2026-08-08T12:00:00Z",
"softwareName": "E:D Market Connector",
"softwareVersion": "5.11.3",
"uploaderID": "abc123"
}},
"message": {}
}}"#,
schema_ref, message,
)
}
fn envelope(schema_ref: &str, message: &str) -> Envelope {
serde_json::from_str(&envelope_json(schema_ref, message))
.expect("envelope should parse")
}
pub(super) fn frame(schema_ref: &str, message: &str) -> Vec<u8> {
miniz_oxide::deflate::compress_to_vec_zlib(
envelope_json(schema_ref, message).as_bytes(),
6,
)
}
pub(super) const JUMP: &str = r#"{
"timestamp": "2026-08-08T12:00:00Z",
"event": "FSDJump",
"StarSystem": "Sol",
"StarPos": [0.0, 0.0, 0.0],
"SystemAddress": 10477373803
}"#;
#[test]
fn a_journal_message_is_a_journal_entry() {
let envelope = envelope("https://eddn.edcd.io/schemas/journal/1", JUMP);
assert!(matches!(envelope.message, Message::Journal(_)));
}
#[test]
fn the_standalone_journal_schemas_read_as_journal_entries() {
for name in super::JOURNAL_SCHEMAS {
let reference = format!("https://eddn.edcd.io/schemas/{}/1", name);
let envelope = envelope(&reference, JUMP);
assert!(
matches!(envelope.message, Message::Journal(_)),
"{} did not read as a journal entry",
name,
);
}
}
#[test]
fn the_schemas_with_no_event_are_placed_by_their_reference() {
let outfitting = envelope(
"https://eddn.edcd.io/schemas/outfitting/3",
r#"{
"timestamp": "2026-08-08T12:00:00Z",
"systemName": "Sol",
"stationName": "Abraham Lincoln",
"marketId": 128016384,
"modules": ["Int_Engine_Size3_Class5_Fast"]
}"#,
);
assert!(matches!(outfitting.message, Message::Outfitting(_)));
let shipyard = envelope(
"https://eddn.edcd.io/schemas/shipyard/2",
r#"{
"timestamp": "2026-08-08T12:00:00Z",
"systemName": "Sol",
"stationName": "Abraham Lincoln",
"marketId": 128016384,
"ships": ["SideWinder"]
}"#,
);
assert!(matches!(shipyard.message, Message::Shipyard(_)));
let black_market = envelope(
"https://eddn.edcd.io/schemas/blackmarket/1",
r#"{
"timestamp": "2026-08-08T12:00:00Z",
"systemName": "Sol",
"stationName": "Abraham Lincoln",
"name": "Gold",
"sellPrice": 9432,
"prohibited": false
}"#,
);
assert!(matches!(black_market.message, Message::BlackMarket(_)));
}
#[test]
fn either_outfitting_version_is_read() {
for (version, modules) in [
("2", r#"["Int_Engine_Size3_Class5_Fast"]"#),
(
"3",
r#"[{
"id": 128064258,
"Name": "Int_Engine_Size3_Class5_Fast",
"BuyPrice": 5103953,
"BuyMercCoinsPrice": 0
}]"#,
),
] {
let reference =
format!("https://eddn.edcd.io/schemas/outfitting/{}", version);
let message = format!(
r#"{{
"timestamp": "2026-08-08T12:00:00Z",
"systemName": "Sol",
"stationName": "Abraham Lincoln",
"marketId": 128016384,
"modules": {}
}}"#,
modules,
);
let envelope = envelope(&reference, &message);
assert!(
matches!(envelope.message, Message::Outfitting(_)),
"outfitting/{} was not read",
version,
);
}
}
#[test]
fn test_schemas_are_read_but_not_live() {
let envelope =
envelope("https://eddn.edcd.io/schemas/journal/1/test", JUMP);
assert!(!envelope.live);
assert!(matches!(envelope.message, Message::Journal(_)));
}
#[test]
fn an_ordinary_schema_is_live() {
assert!(envelope("https://eddn.edcd.io/schemas/journal/1", JUMP).live);
assert!(envelope("nonsense", JUMP).live);
}
#[test]
fn an_unmodelled_schema_is_kept() {
let envelope = envelope(
"https://eddn.edcd.io/schemas/fcmaterials_journal/1",
r#"{
"timestamp": "2026-08-08T12:00:00Z",
"event": "FCMaterials",
"MarketID": 3700571136,
"CarrierID": "K7Q-BQL",
"CarrierName": "Nomad",
"Items": []
}"#,
);
assert!(matches!(envelope.message, Message::Unmodeled(_)));
assert_eq!(
envelope.schema_ref,
"https://eddn.edcd.io/schemas/fcmaterials_journal/1",
);
}
#[test]
fn a_payload_that_will_not_read_is_reported() {
let json = r#"{
"$schemaRef": "https://eddn.edcd.io/schemas/journal/1",
"header": {
"gatewayTimestamp": "2026-08-08T12:00:00Z",
"softwareName": "E:D Market Connector",
"softwareVersion": "5.11.3",
"uploaderID": "abc123"
},
"message": {
"timestamp": "2026-08-08T12:00:00Z",
"event": "FSDJump",
"StarSystem": "Sol",
"StarPos": "nowhere",
"SystemAddress": 10477373803
}
}"#;
assert!(serde_json::from_str::<Envelope>(json).is_err());
}
#[test]
fn a_reference_that_makes_no_sense_is_unmodelled() {
assert_eq!(Schema::read("nonsense"), None);
assert_eq!(Schema::read("https://eddn.edcd.io/schemas/journal"), None);
let envelope = envelope("nonsense", JUMP);
assert!(matches!(envelope.message, Message::Unmodeled(_)));
assert_eq!(envelope.version, None);
}
#[test]
fn a_reference_reads_into_what_it_says() {
assert_eq!(
Schema::read("https://eddn.edcd.io/schemas/commodity/3"),
Some(Schema { name: "commodity", version: "3", live: true }),
);
assert_eq!(
Schema::read("https://eddn.edcd.io/schemas/shipyard/2/test"),
Some(Schema { name: "shipyard", version: "2", live: false }),
);
}
#[test]
fn the_version_is_said_even_though_nothing_turns_on_it() {
let old = envelope(
"https://eddn.edcd.io/schemas/outfitting/2",
r#"{
"timestamp": "2026-08-08T12:00:00Z",
"systemName": "Sol",
"stationName": "Abraham Lincoln",
"marketId": 128016384,
"modules": ["Int_Engine_Size3_Class5_Fast"]
}"#,
);
let new = envelope(
"https://eddn.edcd.io/schemas/outfitting/3",
r#"{
"timestamp": "2026-08-08T12:00:00Z",
"systemName": "Sol",
"stationName": "Abraham Lincoln",
"marketId": 128016384,
"modules": [{
"id": 128064258,
"Name": "Int_Engine_Size3_Class5_Fast",
"BuyPrice": 5103953,
"BuyMercCoinsPrice": 0
}]
}"#,
);
assert_eq!(old.version.as_deref(), Some("2"));
assert_eq!(new.version.as_deref(), Some("3"));
assert!(matches!(old.message, Message::Outfitting(_)));
assert!(matches!(new.message, Message::Outfitting(_)));
}
#[test]
fn a_test_schema_reports_its_version_too() {
let envelope =
envelope("https://eddn.edcd.io/schemas/journal/1/test", JUMP);
assert_eq!(envelope.version.as_deref(), Some("1"));
assert!(!envelope.live);
}
}
#[cfg(test)]
mod frames {
use super::tests::{envelope_json, frame, JUMP};
use super::*;
const JOURNAL: &str = "https://eddn.edcd.io/schemas/journal/1";
const JOURNAL_TEST: &str = "https://eddn.edcd.io/schemas/journal/1/test";
#[test]
fn a_live_frame_is_handed_over() {
let read = read_frame(&frame(JOURNAL, JUMP))
.expect("a live frame should be handed over");
let envelope = read.expect("and should read");
assert!(envelope.live);
assert!(matches!(envelope.message, Message::Journal(_)));
}
#[test]
fn a_test_frame_is_let_go() {
assert!(read_frame(&frame(JOURNAL_TEST, JUMP)).is_none());
}
#[test]
fn a_frame_that_is_not_zlib_is_reported() {
assert!(matches!(
read_frame(b"not zlib at all"),
Some(Err(Error::Decompress(_))),
));
}
#[test]
fn a_frame_that_is_not_an_envelope_is_reported() {
let rubbish = miniz_oxide::deflate::compress_to_vec_zlib(b"{}", 6);
assert!(matches!(read_frame(&rubbish), Some(Err(Error::Parse { .. }))));
}
#[test]
fn a_frame_whose_payload_will_not_read_is_reported() {
let json = envelope_json(
JOURNAL,
r#"{
"timestamp": "2026-08-08T12:00:00Z",
"event": "FSDJump",
"StarSystem": "Sol",
"StarPos": "nowhere",
"SystemAddress": 10477373803
}"#,
);
let bad =
miniz_oxide::deflate::compress_to_vec_zlib(json.as_bytes(), 6);
assert!(matches!(read_frame(&bad), Some(Err(Error::Parse { .. }))));
}
#[test]
fn a_payload_error_keeps_the_message_it_could_not_read() {
let json = envelope_json(
"https://eddn.edcd.io/schemas/fssdiscoveryscan/1",
r#"{
"timestamp": "2026-08-08T12:00:00Z",
"event": "FSSDiscoveryScan",
"SystemName": "Sol",
"StarPos": [0.0, 0.0, 0.0],
"SystemAddress": 10477373803,
"BodyCount": "",
"NonBodyCount": 3,
"Progress": 1.0
}"#,
);
let bad =
miniz_oxide::deflate::compress_to_vec_zlib(json.as_bytes(), 6);
let Some(Err(err)) = read_frame(&bad) else {
panic!("a payload that will not read should be reported")
};
let said = err.to_string();
assert!(said.contains("BodyCount"), "did not name the field: {}", said,);
assert!(err.near().is_none(), "pointed somewhere anyway: {}", said);
}
#[test]
fn a_malformed_envelope_shows_where_it_stopped() {
let truncated = miniz_oxide::deflate::compress_to_vec_zlib(
br#"{"$schemaRef": "https://eddn.edcd.io/schemas/journal/1", "heade"#,
6,
);
let Some(Err(err)) = read_frame(&truncated) else {
panic!("a malformed envelope should be reported")
};
let near = err.near().expect("should point at where it stopped");
assert!(near.contains("heade"), "pointed elsewhere: {}", near);
}
}