///|
pub(all) enum QuerySort {
QueryOriginal
QueryTimestampAsc
QueryTimestampDesc
QueryTopicAsc
} derive(Eq, @debug.Debug)
///|
pub(all) struct HeaderFilter {
key : String
value : String?
} derive(Eq, @debug.Debug)
///|
pub(all) struct EventQuery {
name : String
topic : TopicPattern?
headers : Array[HeaderFilter]
predicates : Array[EventPredicate]
min_timestamp_ms : Int?
max_timestamp_ms : Int?
offset : Int
limit : Int
sort : QuerySort
} derive(Eq, @debug.Debug)
///|
pub(all) struct QuerySkip {
event_id : String
topic : String
reason : String
} derive(Eq, @debug.Debug)
///|
pub(all) struct QueryRow {
index : Int
event : Envelope
fingerprint : String
} derive(Eq, @debug.Debug)
///|
pub(all) struct QueryResult {
name : String
scanned : Int
matched_before_page : Int
returned : Int
skipped : Array[QuerySkip]
rows : Array[QueryRow]
} derive(Eq, @debug.Debug)
///|
pub fn header_filter(
key : StringView,
value? : StringView = "",
) -> HeaderFilter {
let actual = if value == "" { None } else { Some(value.to_owned()) }
{ key: key.to_owned(), value: actual }
}
///|
pub fn event_query(
name? : StringView = "query",
topic? : StringView = "",
headers? : ArrayView[HeaderFilter] = [],
predicates? : ArrayView[EventPredicate] = [],
min_timestamp_ms? : Int,
max_timestamp_ms? : Int,
offset? : Int = 0,
limit? : Int = 0,
sort? : QuerySort = QueryOriginal,
) -> Result[EventQuery, EventRailError] {
let parsed = if topic == "" {
None
} else {
match topic_pattern(topic) {
Ok(value) => Some(value)
Err(err) => return Err(err)
}
}
Ok({
name: name.to_owned(),
topic: parsed,
headers: headers.to_owned(),
predicates: predicates.to_owned(),
min_timestamp_ms,
max_timestamp_ms,
offset: if offset < 0 {
0
} else {
offset
},
limit: if limit < 0 {
0
} else {
limit
},
sort,
})
}
///|
pub fn EventQuery::with_topic(
self : EventQuery,
topic : StringView,
) -> Result[EventQuery, EventRailError] {
match topic_pattern(topic) {
Ok(value) => Ok({ ..self, topic: Some(value) })
Err(err) => Err(err)
}
}
///|
pub fn EventQuery::without_topic(self : EventQuery) -> EventQuery {
{ ..self, topic: None }
}
///|
pub fn EventQuery::with_header(
self : EventQuery,
key : StringView,
value? : StringView = "",
) -> EventQuery {
let headers = self.headers.copy()
headers.push(header_filter(key, value~))
{ ..self, headers, }
}
///|
pub fn EventQuery::with_predicate(
self : EventQuery,
predicate : EventPredicate,
) -> EventQuery {
let predicates = self.predicates.copy()
predicates.push(predicate)
{ ..self, predicates, }
}
///|
pub fn EventQuery::with_time_range(
self : EventQuery,
min_timestamp_ms? : Int,
max_timestamp_ms? : Int,
) -> EventQuery {
{ ..self, min_timestamp_ms, max_timestamp_ms }
}
///|
pub fn EventQuery::with_page(
self : EventQuery,
offset : Int,
limit : Int,
) -> EventQuery {
{
..self,
offset: if offset < 0 {
0
} else {
offset
},
limit: if limit < 0 {
0
} else {
limit
},
}
}
///|
pub fn EventQuery::with_sort(self : EventQuery, sort : QuerySort) -> EventQuery {
{ ..self, sort, }
}
///|
pub fn EventQuery::matches(
self : EventQuery,
event : Envelope,
) -> Result[(Bool, String), EventRailError] {
match self.topic {
Some(pattern) =>
match pattern.matches_topic(event.topic) {
Err(err) => return Err(err)
Ok(false) => return Ok((false, "topic mismatch"))
Ok(true) => ()
}
None => ()
}
match self.min_timestamp_ms {
Some(minimum) if event.timestamp_ms < minimum =>
return Ok((false, "timestamp before minimum"))
_ => ()
}
match self.max_timestamp_ms {
Some(maximum) if event.timestamp_ms > maximum =>
return Ok((false, "timestamp after maximum"))
_ => ()
}
for filter in self.headers {
match event.header(filter.key) {
None => return Ok((false, "missing header \{filter.key}"))
Some(actual) =>
match filter.value {
Some(expected) if actual != expected =>
return Ok((false, "header \{filter.key} mismatch"))
_ => ()
}
}
}
for predicate in self.predicates {
let result = predicate.evaluate(event)
if !result.passed {
return Ok((false, result.reason))
}
}
Ok((true, "matched"))
}
///|
pub fn EventBatch::query(
self : EventBatch,
query : EventQuery,
) -> Result[QueryResult, EventRailError] {
query.apply(self.events)
}
///|
pub fn EventTape::query(
self : EventTape,
query : EventQuery,
) -> Result[QueryResult, EventRailError] {
query.apply(self.entries.map(entry => entry.event))
}
///|
pub fn EventQuery::apply(
self : EventQuery,
events : ArrayView[Envelope],
) -> Result[QueryResult, EventRailError] {
let matched : Array[QueryRow] = []
let skipped : Array[QuerySkip] = []
for index, event in events {
match self.matches(event) {
Err(err) => return Err(err)
Ok((true, _)) =>
matched.push({ index, event, fingerprint: event.fingerprint() })
Ok((false, reason)) =>
skipped.push({ event_id: event.id, topic: event.topic, reason })
}
}
matched.sort_by(fn(left, right) { compare_query_row(left, right, self.sort) })
let rows = page_query_rows(matched, self.offset, self.limit)
Ok({
name: self.name,
scanned: events.length(),
matched_before_page: matched.length(),
returned: rows.length(),
skipped,
rows,
})
}
///|
pub fn QueryResult::events(self : QueryResult) -> Array[Envelope] {
self.rows.map(row => row.event)
}
///|
pub fn QueryResult::is_empty(self : QueryResult) -> Bool {
self.rows.length() == 0
}
///|
pub fn QueryResult::summary(self : QueryResult) -> String {
"query=\{escape_wire_text(self.name)} scanned=\{self.scanned} matched=\{self.matched_before_page} returned=\{self.returned} skipped=\{self.skipped.length()}"
}
///|
pub fn QueryResult::manifest_lines(self : QueryResult) -> Array[String] {
let lines : Array[String] = []
lines.push(self.summary())
for row in self.rows {
lines.push(row.to_wire())
}
lines
}
///|
pub fn QueryResult::manifest(self : QueryResult) -> String {
self.manifest_lines().join("\n")
}
///|
pub fn QueryResult::skip_lines(self : QueryResult) -> Array[String] {
self.skipped.map(skip => skip.to_wire())
}
///|
pub fn QueryRow::to_wire(self : QueryRow) -> String {
"idx=\{self.index};event=\{escape_wire_text(self.event.id)};topic=\{escape_wire_text(self.event.topic)};timestamp=\{self.event.timestamp_ms};fingerprint=\{escape_wire_text(self.fingerprint)}"
}
///|
pub fn QuerySkip::to_wire(self : QuerySkip) -> String {
"event=\{escape_wire_text(self.event_id)};topic=\{escape_wire_text(self.topic)};reason=\{escape_wire_text(self.reason)}"
}
///|
pub fn HeaderFilter::to_wire(self : HeaderFilter) -> String {
match self.value {
Some(value) =>
"header=\{escape_wire_text(self.key)}:\{escape_wire_text(value)}"
None => "header=\{escape_wire_text(self.key)}"
}
}
///|
pub fn EventQuery::to_wire(self : EventQuery) -> String {
let topic = match self.topic {
Some(pattern) => pattern.raw
None => "*"
}
"query=\{escape_wire_text(self.name)};topic=\{escape_wire_text(topic)};headers=\{self.headers.map(h => h.to_wire()).join(",")};predicates=\{self.predicates.length()};offset=\{self.offset};limit=\{self.limit};sort=\{self.sort.to_wire()}"
}
///|
pub fn QuerySort::to_wire(self : QuerySort) -> String {
match self {
QueryOriginal => "original"
QueryTimestampAsc => "timestamp-asc"
QueryTimestampDesc => "timestamp-desc"
QueryTopicAsc => "topic-asc"
}
}
///|
fn page_query_rows(
rows : Array[QueryRow],
offset : Int,
limit : Int,
) -> Array[QueryRow] {
let paged : Array[QueryRow] = []
for index, row in rows {
if index < offset {
continue
}
if limit > 0 && paged.length() >= limit {
break
}
paged.push(row)
}
paged
}
///|
fn compare_query_row(
left : QueryRow,
right : QueryRow,
sort : QuerySort,
) -> Int {
match sort {
QueryOriginal => left.index - right.index
QueryTimestampAsc => compare_query_timestamp(left, right)
QueryTimestampDesc => compare_query_timestamp(right, left)
QueryTopicAsc =>
if left.event.topic != right.event.topic {
left.event.topic.compare(right.event.topic)
} else {
compare_query_timestamp(left, right)
}
}
}
///|
fn compare_query_timestamp(left : QueryRow, right : QueryRow) -> Int {
if left.event.timestamp_ms != right.event.timestamp_ms {
left.event.timestamp_ms - right.event.timestamp_ms
} else if left.event.topic != right.event.topic {
left.event.topic.compare(right.event.topic)
} else {
left.event.id.compare(right.event.id)
}
}