487 lines
16 KiB
GDScript
487 lines
16 KiB
GDScript
class_name WorldEventRetentionPlan
|
|
extends RefCounted
|
|
|
|
const SCHEMA_VERSION := 1
|
|
const DECISION_RETAIN_RAW := "retain_raw"
|
|
const DECISION_ROLLED_UP := "rolled_up"
|
|
const RECORD_FIELDS := [
|
|
"schema_version",
|
|
"current_tick",
|
|
"ticks_per_day",
|
|
"raw_retention_days",
|
|
"first_raw_day_index",
|
|
"source_checksum",
|
|
"high_volume_event_types",
|
|
"pinned_fact_ids",
|
|
"causal_fact_ids",
|
|
"retained_event_ids",
|
|
"rolled_up_event_ids",
|
|
"rollup_store",
|
|
"reason_traces",
|
|
"telemetry",
|
|
]
|
|
const TRACE_FIELDS := [
|
|
"event_id",
|
|
"decision",
|
|
"reasons",
|
|
"cause_event_ids",
|
|
"required_by_event_ids",
|
|
"rollup_keys",
|
|
]
|
|
const TELEMETRY_FIELDS := [
|
|
"succeeded",
|
|
"input_event_count",
|
|
"retained_event_count",
|
|
"rolled_up_event_count",
|
|
"rollup_record_count",
|
|
"rollup_contribution_count",
|
|
"pinned_fact_count",
|
|
"causal_fact_count",
|
|
"causal_ancestor_count",
|
|
"recent_event_count",
|
|
"non_rollup_event_count",
|
|
"invalid_rollup_event_count",
|
|
"first_raw_day_index",
|
|
"source_checksum",
|
|
"source_unchanged",
|
|
]
|
|
|
|
var _data: Dictionary = {}
|
|
var _rollup_store := WorldEventDailyRollupStore.new()
|
|
var _traces_by_event_id: Dictionary = {}
|
|
|
|
|
|
static func create(
|
|
current_tick: int,
|
|
ticks_per_day: int,
|
|
raw_retention_days: int,
|
|
first_raw_day_index: int,
|
|
source_checksum: String,
|
|
high_volume_event_types: Array[StringName],
|
|
pinned_fact_ids: Array[int],
|
|
causal_fact_ids: Array[int],
|
|
retained_event_ids: Array[int],
|
|
rolled_up_event_ids: Array[int],
|
|
rollup_store: WorldEventDailyRollupStore,
|
|
reason_traces: Array[Dictionary],
|
|
telemetry: Dictionary
|
|
) -> WorldEventRetentionPlan:
|
|
if rollup_store == null:
|
|
return null
|
|
var serialized_event_types: Array[String] = []
|
|
for event_type: StringName in high_volume_event_types:
|
|
serialized_event_types.append(String(event_type))
|
|
return from_dictionary(
|
|
{
|
|
"schema_version": SCHEMA_VERSION,
|
|
"current_tick": current_tick,
|
|
"ticks_per_day": ticks_per_day,
|
|
"raw_retention_days": raw_retention_days,
|
|
"first_raw_day_index": first_raw_day_index,
|
|
"source_checksum": source_checksum,
|
|
"high_volume_event_types": serialized_event_types,
|
|
"pinned_fact_ids": pinned_fact_ids.duplicate(),
|
|
"causal_fact_ids": causal_fact_ids.duplicate(),
|
|
"retained_event_ids": retained_event_ids.duplicate(),
|
|
"rolled_up_event_ids": rolled_up_event_ids.duplicate(),
|
|
"rollup_store": rollup_store.to_dictionary(),
|
|
"reason_traces": reason_traces.duplicate(true),
|
|
"telemetry": telemetry.duplicate(true),
|
|
}
|
|
)
|
|
|
|
|
|
static func from_dictionary(record_data: Dictionary) -> WorldEventRetentionPlan:
|
|
if not _has_exact_fields(record_data, RECORD_FIELDS):
|
|
return null
|
|
if (
|
|
not record_data["schema_version"] is int
|
|
or int(record_data["schema_version"]) != SCHEMA_VERSION
|
|
or not record_data["current_tick"] is int
|
|
or not record_data["ticks_per_day"] is int
|
|
or not record_data["raw_retention_days"] is int
|
|
or not record_data["first_raw_day_index"] is int
|
|
or not record_data["source_checksum"] is String
|
|
or not record_data["high_volume_event_types"] is Array
|
|
or not record_data["pinned_fact_ids"] is Array
|
|
or not record_data["causal_fact_ids"] is Array
|
|
or not record_data["retained_event_ids"] is Array
|
|
or not record_data["rolled_up_event_ids"] is Array
|
|
or not record_data["rollup_store"] is Dictionary
|
|
or not record_data["reason_traces"] is Array
|
|
or not record_data["telemetry"] is Dictionary
|
|
):
|
|
return null
|
|
var event_types: Variant = _normalize_string_array(record_data["high_volume_event_types"])
|
|
var pinned_ids: Variant = _normalize_int_array(record_data["pinned_fact_ids"])
|
|
var causal_ids: Variant = _normalize_int_array(record_data["causal_fact_ids"])
|
|
var retained_ids: Variant = _normalize_int_array(record_data["retained_event_ids"])
|
|
var rolled_up_ids: Variant = _normalize_int_array(record_data["rolled_up_event_ids"])
|
|
if (
|
|
event_types == null
|
|
or pinned_ids == null
|
|
or causal_ids == null
|
|
or retained_ids == null
|
|
or rolled_up_ids == null
|
|
or event_types != record_data["high_volume_event_types"]
|
|
or pinned_ids != record_data["pinned_fact_ids"]
|
|
or causal_ids != record_data["causal_fact_ids"]
|
|
or retained_ids != record_data["retained_event_ids"]
|
|
or rolled_up_ids != record_data["rolled_up_event_ids"]
|
|
):
|
|
return null
|
|
var current_tick := int(record_data["current_tick"])
|
|
var ticks_per_day := int(record_data["ticks_per_day"])
|
|
var raw_retention_days := int(record_data["raw_retention_days"])
|
|
var first_raw_day_index := int(record_data["first_raw_day_index"])
|
|
if (
|
|
current_tick < 0
|
|
or ticks_per_day <= 0
|
|
or raw_retention_days < 0
|
|
or (first_raw_day_index != maxi(int(current_tick / ticks_per_day) - raw_retention_days, 0))
|
|
or not _is_sha256(String(record_data["source_checksum"]))
|
|
):
|
|
return null
|
|
var rollup_store := WorldEventDailyRollupStore.from_dictionary(record_data["rollup_store"])
|
|
if rollup_store == null:
|
|
return null
|
|
var normalized_traces: Variant = _normalize_traces(record_data["reason_traces"])
|
|
if normalized_traces == null or normalized_traces != record_data["reason_traces"]:
|
|
return null
|
|
var telemetry: Dictionary = record_data["telemetry"]
|
|
if not _telemetry_is_valid(telemetry, String(record_data["source_checksum"])):
|
|
return null
|
|
if not _cross_references_are_valid(
|
|
pinned_ids,
|
|
causal_ids,
|
|
retained_ids,
|
|
rolled_up_ids,
|
|
rollup_store,
|
|
normalized_traces,
|
|
telemetry,
|
|
first_raw_day_index
|
|
):
|
|
return null
|
|
var plan := WorldEventRetentionPlan.new()
|
|
plan._data = {
|
|
"current_tick": current_tick,
|
|
"ticks_per_day": ticks_per_day,
|
|
"raw_retention_days": raw_retention_days,
|
|
"first_raw_day_index": first_raw_day_index,
|
|
"source_checksum": String(record_data["source_checksum"]),
|
|
"high_volume_event_types": event_types,
|
|
"pinned_fact_ids": pinned_ids,
|
|
"causal_fact_ids": causal_ids,
|
|
"retained_event_ids": retained_ids,
|
|
"rolled_up_event_ids": rolled_up_ids,
|
|
"reason_traces": normalized_traces,
|
|
"telemetry": telemetry.duplicate(true),
|
|
}
|
|
plan._rollup_store = rollup_store
|
|
for trace: Dictionary in normalized_traces:
|
|
plan._traces_by_event_id[int(trace["event_id"])] = trace.duplicate(true)
|
|
return plan
|
|
|
|
|
|
func get_retained_event_ids() -> Array[int]:
|
|
return _copy_int_array(_data["retained_event_ids"])
|
|
|
|
|
|
func get_rolled_up_event_ids() -> Array[int]:
|
|
return _copy_int_array(_data["rolled_up_event_ids"])
|
|
|
|
|
|
func get_rollup_store() -> WorldEventDailyRollupStore:
|
|
return WorldEventDailyRollupStore.from_dictionary(_rollup_store.to_dictionary())
|
|
|
|
|
|
func get_reason_trace(event_id: int) -> Dictionary:
|
|
var trace: Dictionary = _traces_by_event_id.get(event_id, {} as Dictionary)
|
|
return trace.duplicate(true)
|
|
|
|
|
|
func get_reason_traces() -> Array[Dictionary]:
|
|
var traces: Array[Dictionary] = []
|
|
for trace: Dictionary in _data["reason_traces"]:
|
|
traces.append(trace.duplicate(true))
|
|
return traces
|
|
|
|
|
|
func get_telemetry() -> Dictionary:
|
|
return (_data["telemetry"] as Dictionary).duplicate(true)
|
|
|
|
|
|
func get_source_checksum() -> String:
|
|
return String(_data["source_checksum"])
|
|
|
|
|
|
func to_dictionary() -> Dictionary:
|
|
return {
|
|
"schema_version": SCHEMA_VERSION,
|
|
"current_tick": int(_data["current_tick"]),
|
|
"ticks_per_day": int(_data["ticks_per_day"]),
|
|
"raw_retention_days": int(_data["raw_retention_days"]),
|
|
"first_raw_day_index": int(_data["first_raw_day_index"]),
|
|
"source_checksum": String(_data["source_checksum"]),
|
|
"high_volume_event_types": (_data["high_volume_event_types"] as Array).duplicate(),
|
|
"pinned_fact_ids": (_data["pinned_fact_ids"] as Array).duplicate(),
|
|
"causal_fact_ids": (_data["causal_fact_ids"] as Array).duplicate(),
|
|
"retained_event_ids": (_data["retained_event_ids"] as Array).duplicate(),
|
|
"rolled_up_event_ids": (_data["rolled_up_event_ids"] as Array).duplicate(),
|
|
"rollup_store": _rollup_store.to_dictionary(),
|
|
"reason_traces": (_data["reason_traces"] as Array).duplicate(true),
|
|
"telemetry": (_data["telemetry"] as Dictionary).duplicate(true),
|
|
}
|
|
|
|
|
|
func checksum() -> String:
|
|
return JSON.stringify(to_dictionary()).sha256_text()
|
|
|
|
|
|
static func _normalize_traces(raw_traces: Array) -> Variant:
|
|
var traces: Array[Dictionary] = []
|
|
var seen: Dictionary = {}
|
|
var last_event_id := -1
|
|
for raw_trace: Variant in raw_traces:
|
|
if not raw_trace is Dictionary or not _has_exact_fields(raw_trace, TRACE_FIELDS):
|
|
return null
|
|
if (
|
|
not raw_trace["event_id"] is int
|
|
or not raw_trace["decision"] is String
|
|
or not raw_trace["reasons"] is Array
|
|
or not raw_trace["cause_event_ids"] is Array
|
|
or not raw_trace["required_by_event_ids"] is Array
|
|
or not raw_trace["rollup_keys"] is Array
|
|
):
|
|
return null
|
|
var event_id := int(raw_trace["event_id"])
|
|
var decision := String(raw_trace["decision"])
|
|
var reasons: Variant = _normalize_string_array(raw_trace["reasons"])
|
|
var cause_ids: Variant = _normalize_int_array(raw_trace["cause_event_ids"])
|
|
var required_by_ids: Variant = _normalize_int_array(raw_trace["required_by_event_ids"])
|
|
var rollup_keys: Variant = _normalize_string_array(raw_trace["rollup_keys"])
|
|
if (
|
|
event_id < 0
|
|
or event_id <= last_event_id
|
|
or seen.has(event_id)
|
|
or decision not in [DECISION_RETAIN_RAW, DECISION_ROLLED_UP]
|
|
or reasons == null
|
|
or reasons.is_empty()
|
|
or cause_ids == null
|
|
or required_by_ids == null
|
|
or rollup_keys == null
|
|
or reasons != raw_trace["reasons"]
|
|
or cause_ids != raw_trace["cause_event_ids"]
|
|
or required_by_ids != raw_trace["required_by_event_ids"]
|
|
or rollup_keys != raw_trace["rollup_keys"]
|
|
):
|
|
return null
|
|
seen[event_id] = true
|
|
last_event_id = event_id
|
|
traces.append(raw_trace.duplicate(true))
|
|
return traces
|
|
|
|
|
|
static func _telemetry_is_valid(telemetry: Dictionary, source_checksum: String) -> bool:
|
|
if not _has_exact_fields(telemetry, TELEMETRY_FIELDS):
|
|
return false
|
|
if (
|
|
telemetry["succeeded"] is not bool
|
|
or not bool(telemetry["succeeded"])
|
|
or telemetry["source_unchanged"] is not bool
|
|
or not bool(telemetry["source_unchanged"])
|
|
or telemetry["source_checksum"] is not String
|
|
or String(telemetry["source_checksum"]) != source_checksum
|
|
):
|
|
return false
|
|
for field: String in TELEMETRY_FIELDS:
|
|
if field in ["succeeded", "source_unchanged", "source_checksum"]:
|
|
continue
|
|
if telemetry[field] is not int or int(telemetry[field]) < 0:
|
|
return false
|
|
return true
|
|
|
|
|
|
static func _cross_references_are_valid(
|
|
pinned_ids: Array,
|
|
causal_ids: Array,
|
|
retained_ids: Array,
|
|
rolled_up_ids: Array,
|
|
rollup_store: WorldEventDailyRollupStore,
|
|
traces: Array,
|
|
telemetry: Dictionary,
|
|
first_raw_day_index: int
|
|
) -> bool:
|
|
var retained_set := _id_set(retained_ids)
|
|
var rolled_up_set := _id_set(rolled_up_ids)
|
|
var traces_by_id: Dictionary = {}
|
|
for trace: Dictionary in traces:
|
|
traces_by_id[int(trace["event_id"])] = trace
|
|
for event_id: int in retained_ids:
|
|
if rolled_up_set.has(event_id):
|
|
return false
|
|
for event_id: int in pinned_ids:
|
|
if (
|
|
not retained_set.has(event_id)
|
|
or not traces_by_id.has(event_id)
|
|
or "pinned_fact" not in (traces_by_id[event_id] as Dictionary)["reasons"]
|
|
):
|
|
return false
|
|
for event_id: int in causal_ids:
|
|
if (
|
|
not retained_set.has(event_id)
|
|
or not traces_by_id.has(event_id)
|
|
or "explicit_causal_fact" not in (traces_by_id[event_id] as Dictionary)["reasons"]
|
|
):
|
|
return false
|
|
var rollup_keys: Dictionary = {}
|
|
var contribution_count := 0
|
|
for record: WorldEventDailyRollupRecord in rollup_store.get_all():
|
|
rollup_keys[record.index_key()] = true
|
|
contribution_count += record.get_event_count()
|
|
var trace_ids: Dictionary = {}
|
|
var causal_ancestor_count := 0
|
|
var recent_count := 0
|
|
var non_rollup_count := 0
|
|
var invalid_rollup_count := 0
|
|
var source_ids_by_rollup_key: Dictionary = {}
|
|
for trace: Dictionary in traces:
|
|
var event_id := int(trace["event_id"])
|
|
trace_ids[event_id] = true
|
|
var is_retained := retained_set.has(event_id)
|
|
var is_rolled_up := rolled_up_set.has(event_id)
|
|
if is_retained == is_rolled_up:
|
|
return false
|
|
if is_retained and String(trace["decision"]) != DECISION_RETAIN_RAW:
|
|
return false
|
|
if is_rolled_up and String(trace["decision"]) != DECISION_ROLLED_UP:
|
|
return false
|
|
var reasons: Array = trace["reasons"]
|
|
var trace_rollup_keys: Array = trace["rollup_keys"]
|
|
if is_retained and not trace_rollup_keys.is_empty():
|
|
return false
|
|
if is_retained and "rolled_up_daily" in reasons:
|
|
return false
|
|
if is_rolled_up and (reasons != ["rolled_up_daily"] or trace_rollup_keys.is_empty()):
|
|
return false
|
|
for rollup_key: String in trace_rollup_keys:
|
|
if not rollup_keys.has(rollup_key):
|
|
return false
|
|
if not source_ids_by_rollup_key.has(rollup_key):
|
|
source_ids_by_rollup_key[rollup_key] = []
|
|
var source_ids: Array = source_ids_by_rollup_key[rollup_key]
|
|
source_ids.append(event_id)
|
|
for cause_id: int in trace["cause_event_ids"]:
|
|
if not retained_set.has(cause_id) or not traces_by_id.has(cause_id):
|
|
return false
|
|
var cause_trace: Dictionary = traces_by_id[cause_id]
|
|
if (
|
|
"causal_ancestor" not in cause_trace["reasons"]
|
|
or event_id not in cause_trace["required_by_event_ids"]
|
|
):
|
|
return false
|
|
for required_by_id: int in trace["required_by_event_ids"]:
|
|
if not traces_by_id.has(required_by_id):
|
|
return false
|
|
var required_by_trace: Dictionary = traces_by_id[required_by_id]
|
|
if event_id not in required_by_trace["cause_event_ids"]:
|
|
return false
|
|
if "causal_ancestor" in reasons:
|
|
causal_ancestor_count += 1
|
|
if "within_raw_retention" in reasons:
|
|
recent_count += 1
|
|
if "not_rollup_eligible" in reasons:
|
|
non_rollup_count += 1
|
|
if "rollup_payload_invalid" in reasons:
|
|
invalid_rollup_count += 1
|
|
if trace_ids.size() != retained_ids.size() + rolled_up_ids.size():
|
|
return false
|
|
if source_ids_by_rollup_key.size() != rollup_store.size():
|
|
return false
|
|
for record: WorldEventDailyRollupRecord in rollup_store.get_all():
|
|
var rollup_key := record.index_key()
|
|
if not source_ids_by_rollup_key.has(rollup_key):
|
|
return false
|
|
var source_ids: Array = source_ids_by_rollup_key[rollup_key]
|
|
source_ids.sort()
|
|
if (
|
|
record.get_event_count() != source_ids.size()
|
|
or record.get_minimum_event_id() != int(source_ids.front())
|
|
or record.get_maximum_event_id() != int(source_ids.back())
|
|
or record.get_source_event_ids_checksum() != JSON.stringify(source_ids).sha256_text()
|
|
):
|
|
return false
|
|
return (
|
|
int(telemetry["input_event_count"]) == trace_ids.size()
|
|
and int(telemetry["retained_event_count"]) == retained_ids.size()
|
|
and int(telemetry["rolled_up_event_count"]) == rolled_up_ids.size()
|
|
and int(telemetry["rollup_record_count"]) == rollup_store.size()
|
|
and int(telemetry["rollup_contribution_count"]) == contribution_count
|
|
and int(telemetry["pinned_fact_count"]) == pinned_ids.size()
|
|
and int(telemetry["causal_fact_count"]) == causal_ids.size()
|
|
and int(telemetry["causal_ancestor_count"]) == causal_ancestor_count
|
|
and int(telemetry["recent_event_count"]) == recent_count
|
|
and int(telemetry["non_rollup_event_count"]) == non_rollup_count
|
|
and int(telemetry["invalid_rollup_event_count"]) == invalid_rollup_count
|
|
and int(telemetry["first_raw_day_index"]) == first_raw_day_index
|
|
)
|
|
|
|
|
|
static func _normalize_string_array(values: Array) -> Variant:
|
|
var seen: Dictionary = {}
|
|
var normalized: Array[String] = []
|
|
for value: Variant in values:
|
|
if (not value is String and not value is StringName) or String(value).is_empty():
|
|
return null
|
|
var text := String(value)
|
|
if seen.has(text):
|
|
return null
|
|
seen[text] = true
|
|
normalized.append(text)
|
|
normalized.sort()
|
|
return normalized
|
|
|
|
|
|
static func _normalize_int_array(values: Array) -> Variant:
|
|
var seen: Dictionary = {}
|
|
var normalized: Array[int] = []
|
|
for value: Variant in values:
|
|
if not value is int or int(value) < 0 or seen.has(int(value)):
|
|
return null
|
|
seen[int(value)] = true
|
|
normalized.append(int(value))
|
|
normalized.sort()
|
|
return normalized
|
|
|
|
|
|
static func _id_set(values: Array) -> Dictionary:
|
|
var result: Dictionary = {}
|
|
for value: int in values:
|
|
result[value] = true
|
|
return result
|
|
|
|
|
|
static func _copy_int_array(values: Array) -> Array[int]:
|
|
var copied: Array[int] = []
|
|
for value: int in values:
|
|
copied.append(value)
|
|
return copied
|
|
|
|
|
|
static func _has_exact_fields(record_data: Dictionary, fields: Array) -> bool:
|
|
if record_data.size() != fields.size():
|
|
return false
|
|
for field: String in fields:
|
|
if not record_data.has(field):
|
|
return false
|
|
return true
|
|
|
|
|
|
static func _is_sha256(value: String) -> bool:
|
|
if value.length() != 64:
|
|
return false
|
|
for character_index in range(value.length()):
|
|
if value[character_index] not in "0123456789abcdef":
|
|
return false
|
|
return true
|