diff --git a/README.md b/README.md index 10f41a1..68523ed 100644 --- a/README.md +++ b/README.md @@ -20,7 +20,9 @@ Copy `addon/` into `res://addons/@aviorstudio_gd-telemetry/` and enable the plug const TelemetryModule = preload("res://addons/@aviorstudio_gd-telemetry/src/telemetry_module.gd") var telemetry := TelemetryModule.new() -telemetry.configure(TelemetryModule.TelemetryConfig.new(true, 50, 0.5)) +telemetry.configure(TelemetryModule.TelemetryConfig.new( + true, 50, 0.5, 50, Callable(self, "_send_batch_to_backend") +)) telemetry.add_event(telemetry.build_event( Time.get_ticks_msec(), @@ -32,8 +34,13 @@ telemetry.add_event(telemetry.build_event( )) if telemetry.should_flush(): - var batch: Array = telemetry.drain_serialized_batch() - _send_batch_to_backend(batch) + var outcome := await telemetry.flush() + if outcome != TelemetryModule.FlushResult.ACKNOWLEDGED: + push_warning("Telemetry remains queued: %s" % outcome) + +func _send_batch_to_backend(batch: Array[Dictionary]) -> bool: + # Return true only after the transport acknowledges this batch. + return await transport.send(batch) ``` ## Event Shape @@ -54,7 +61,29 @@ Serialized events use this dictionary shape: - `build_event`: create consistent events. - `add_event`: queue events. - `should_flush`: check batch size/time thresholds. -- `drain_serialized_batch` / `flush`: hand events to your transport layer. +- `flush`: deliver one FIFO batch and remove it only after a boolean `true` + callback acknowledgement. +- Explicit `AddResult`, `FlushResult`, and `counters()` diagnostics. + +## Delivery And Backpressure Contract + +- Memory only; this addon does not provide durable delivery. +- The maximum and default cap is 1,000 events across queued and in-flight data. At + capacity, the oldest queued event is dropped for a new event. An in-flight + batch is never changed; if it alone consumes capacity, the newest event is + rejected. Both cases have separate counters and `add_event` results. +- One batch may be in flight. Concurrent `flush` calls return `BUSY`. +- A callback returns or asynchronously resolves to boolean `true` to + acknowledge. `false`, an invalid return, or a 10 second timeout fails an + attempt. Defaults are three retries after the initial attempt with capped + exponential delays of 0.5, 1, and 2 seconds. +- A permanently unserializable head event remains queued and reports + `SERIALIZATION_FAILED`; the caller may explicitly discard it with + `discard_oldest_event()`. +- `shutdown(true)` stops the owned timer, attempts a final flush, and returns + the observable outcome. If another flush is pending, shutdown cancels it, + restores its batch, and returns `CANCELLED`. No new events are accepted + afterward. ## Notes diff --git a/addon/plugin.cfg b/addon/plugin.cfg index 47907f4..db59129 100644 --- a/addon/plugin.cfg +++ b/addon/plugin.cfg @@ -2,5 +2,5 @@ name="GD Telemetry" description="Lightweight, game-agnostic telemetry event batching." author="Avior Studio" -version="0.0.1" +version="0.0.2" script="plugin.gd" diff --git a/addon/src/telemetry_module.gd b/addon/src/telemetry_module.gd index 05ff7c3..a844572 100644 --- a/addon/src/telemetry_module.gd +++ b/addon/src/telemetry_module.gd @@ -1,30 +1,56 @@ class_name TelemetryModule extends RefCounted -## Telemetry batching module with configurable flush behavior. +## Bounded, acknowledged, in-memory telemetry delivery. + +signal _ack_completed(generation: int, acknowledged: bool) +signal _flush_requested +signal flush_finished(result) + +enum AddResult { ADDED, DISABLED, INVALID_EVENT, DROPPED_OLDEST_AND_ADDED, REJECTED_IN_FLIGHT_CAPACITY } +enum FlushResult { ACKNOWLEDGED, EMPTY, DISABLED, INVALID_CALLBACK, BUSY, SERIALIZATION_FAILED, FAILED_AFTER_RETRIES, CANCELLED } + +const MAX_QUEUE_EVENTS := 1000 +const MAX_RETRIES := 3 +const MAX_RETRY_INITIAL_DELAY_S := 0.5 +const MAX_RETRY_DELAY_S := 2.0 +const MAX_CALLBACK_TIMEOUT_S := 10.0 class TelemetryConfig extends RefCounted: - ## Telemetry runtime configuration. var enabled: bool var batch_size: int var batch_interval_s: float var max_debug_messages: int var flush_callback: Callable + var max_queue_events: int + var max_retries: int + var retry_initial_delay_s: float + var retry_max_delay_s: float + var callback_timeout_s: float func _init( enabled: bool = false, batch_size: int = 50, batch_interval_s: float = 0.5, max_debug_messages: int = 50, - flush_callback: Callable = Callable() + flush_callback: Callable = Callable(), + max_queue_events: int = 1000, + max_retries: int = 3, + retry_initial_delay_s: float = 0.5, + retry_max_delay_s: float = 2.0, + callback_timeout_s: float = 10.0 ) -> void: self.enabled = enabled self.batch_size = batch_size self.batch_interval_s = batch_interval_s self.max_debug_messages = max_debug_messages self.flush_callback = flush_callback + self.max_queue_events = max_queue_events + self.max_retries = max_retries + self.retry_initial_delay_s = retry_initial_delay_s + self.retry_max_delay_s = retry_max_delay_s + self.callback_timeout_s = callback_timeout_s class TelemetryEvent extends RefCounted: - ## Typed telemetry event payload. var timestamp_msec: int var level: String var context_id: String @@ -32,14 +58,7 @@ class TelemetryEvent extends RefCounted: var message: String var metadata: Dictionary - func _init( - timestamp_msec: int = 0, - level: String = "", - context_id: String = "", - subject_id: String = "", - message: String = "", - metadata: Dictionary = {} - ) -> void: + func _init(timestamp_msec: int = 0, level: String = "", context_id: String = "", subject_id: String = "", message: String = "", metadata: Dictionary = {}) -> void: self.timestamp_msec = timestamp_msec self.level = level self.context_id = context_id @@ -48,116 +67,275 @@ class TelemetryEvent extends RefCounted: self.metadata = metadata.duplicate(true) var _config: TelemetryConfig = TelemetryConfig.new() -var _event_batch: Array[TelemetryEvent] = [] +var _event_queue: Array[TelemetryEvent] = [] +var _in_flight: Array[TelemetryEvent] = [] +var _flushing := false +var _cancel_requested := false +var _stopped := false var _auto_flush_timer: Timer = null +var _auto_flush_owner: Node = null +var _ack_generation := 0 +var _ack_pending_generation := 0 +var _counters := { + "acknowledged_batches": 0, + "callback_failures": 0, + "callback_timeouts": 0, + "dropped_oldest": 0, + "dropped_newest": 0, + "retry_attempts": 0, + "serialization_failures": 0, +} -## Applies telemetry runtime configuration. -func configure(config: TelemetryConfig) -> void: - _config = config if config else TelemetryConfig.new() +func _init() -> void: + _flush_requested.connect(Callable(self, "flush")) + +## Rejects invalid limits without replacing the active configuration. +func configure(config: TelemetryConfig) -> bool: + if config == null or config.batch_size <= 0 or config.max_queue_events <= 0 or config.max_queue_events > MAX_QUEUE_EVENTS: + return false + if config.batch_size > config.max_queue_events or config.batch_interval_s <= 0.0 or event_count() > config.max_queue_events: + return false + if config.max_retries < 0 or config.max_retries > MAX_RETRIES or config.retry_initial_delay_s < 0.0 or config.retry_initial_delay_s > MAX_RETRY_INITIAL_DELAY_S: + return false + if config.retry_max_delay_s < config.retry_initial_delay_s or config.retry_max_delay_s > MAX_RETRY_DELAY_S or config.callback_timeout_s <= 0.0 or config.callback_timeout_s > MAX_CALLBACK_TIMEOUT_S: + return false + _config = config + _stopped = false if _auto_flush_timer and is_instance_valid(_auto_flush_timer): - _auto_flush_timer.wait_time = maxf(_config.batch_interval_s, 0.01) + _auto_flush_timer.wait_time = _config.batch_interval_s + if _config.enabled: + _auto_flush_timer.start() + else: + _auto_flush_timer.stop() + return true -## Returns whether telemetry collection is enabled. func is_enabled() -> bool: return _config.enabled -## Builds a typed telemetry event payload. -func build_event( - timestamp_msec: int, - level: String, - context_id: String, - subject_id: String, - message: String, - metadata: Dictionary -) -> TelemetryEvent: +func build_event(timestamp_msec: int, level: String, context_id: String, subject_id: String, message: String, metadata: Dictionary) -> TelemetryEvent: return TelemetryEvent.new(timestamp_msec, level, context_id, subject_id, message, metadata) -## Adds one event to the current batch and flushes when thresholds are met. -func add_event(event: TelemetryEvent) -> void: - if not _config.enabled or event == null: - return - _event_batch.append(event) - if should_flush() and _config.flush_callback.is_valid(): - flush() +## Queue capacity includes queued and in-flight events. In-flight data is immutable. +func add_event(event: TelemetryEvent) -> AddResult: + if not _config.enabled or _stopped: + return AddResult.DISABLED + if event == null: + return AddResult.INVALID_EVENT + var result := AddResult.ADDED + if event_count() >= _config.max_queue_events: + if not _event_queue.is_empty(): + _event_queue.pop_front() + _counters.dropped_oldest += 1 + result = AddResult.DROPPED_OLDEST_AND_ADDED + else: + _counters.dropped_newest += 1 + return AddResult.REJECTED_IN_FLIGHT_CAPACITY + _event_queue.append(event) + if should_flush() and _config.flush_callback.is_valid() and not _flushing: + _flush_requested.emit() + return result -## Returns true when the in-memory batch has reached flush size. func should_flush() -> bool: - if not _config.enabled: - return false - if _config.batch_size <= 0: - return false - return _event_batch.size() >= _config.batch_size + return _config.enabled and not _stopped and _event_queue.size() >= _config.batch_size -## Returns the current number of queued events. func event_count() -> int: - return _event_batch.size() + return _event_queue.size() + _in_flight.size() -## Serializes and drains the current batch. -func drain_serialized_batch() -> Array[Dictionary]: - if not _config.enabled or _event_batch.is_empty(): - _event_batch.clear() - return [] +func queued_count() -> int: + return _event_queue.size() - var serialized_events: Array[Dictionary] = [] - for event in _event_batch: - serialized_events.append(to_dict(event)) +func in_flight_count() -> int: + return _in_flight.size() - _event_batch.clear() - return serialized_events +func counters() -> Dictionary: + return _counters.duplicate(true) + +## Explicit recovery for a permanently unserializable head event. +func discard_oldest_event() -> bool: + if _event_queue.is_empty(): + return false + _event_queue.pop_front() + _counters.dropped_oldest += 1 + return true -## Flushes the current batch through the configured callback. -func flush() -> void: - if not _config.enabled: - return +## Flushes one FIFO batch. Only a true callback acknowledgement removes it. +func flush() -> FlushResult: + if _flushing: + return FlushResult.BUSY + if not _config.enabled or _stopped: + return FlushResult.DISABLED + if _event_queue.is_empty(): + return FlushResult.EMPTY if not _config.flush_callback.is_valid(): - return - var serialized_batch: Array[Dictionary] = drain_serialized_batch() - if serialized_batch.is_empty(): - return - _config.flush_callback.call(serialized_batch) - -## Starts timer-based automatic flush using a provided owner node. -func start_auto_flush(owner: Node) -> void: - if owner == null: - return + return FlushResult.INVALID_CALLBACK + + var take := mini(_config.batch_size, _event_queue.size()) + for index in range(take): + _in_flight.append(_event_queue[index]) + _event_queue = _event_queue.slice(take) + var serialized := _serialize_in_flight() + if serialized.is_empty() and not _in_flight.is_empty(): + _restore_in_flight() + _counters.serialization_failures += 1 + return FlushResult.SERIALIZATION_FAILED + + _flushing = true + _cancel_requested = false + for attempt in range(_config.max_retries + 1): + if attempt > 0: + _counters.retry_attempts += 1 + var delay := minf(_config.retry_initial_delay_s * pow(2.0, attempt - 1), _config.retry_max_delay_s) + if delay > 0.0: + await Engine.get_main_loop().create_timer(delay).timeout + var acknowledged := await _call_with_timeout(serialized) + if _cancel_requested: + _restore_in_flight() + _flushing = false + flush_finished.emit(FlushResult.CANCELLED) + return FlushResult.CANCELLED + if acknowledged: + _in_flight.clear() + _flushing = false + _counters.acknowledged_batches += 1 + flush_finished.emit(FlushResult.ACKNOWLEDGED) + return FlushResult.ACKNOWLEDGED + _counters.callback_failures += 1 + _restore_in_flight() + _flushing = false + flush_finished.emit(FlushResult.FAILED_AFTER_RETRIES) + return FlushResult.FAILED_AFTER_RETRIES + +func shutdown(flush_pending: bool = true) -> FlushResult: stop_auto_flush() + if _flushing: + call_deferred("cancel_pending_flush") + var cancelled: FlushResult = await flush_finished + _stopped = true + return cancelled + if flush_pending and _config.enabled and event_count() > 0: + var result: FlushResult = await flush() + _stopped = true + return result + _stopped = true + return FlushResult.EMPTY + +func cancel_pending_flush() -> bool: + if not _flushing: + return false + _cancel_requested = true + if _ack_pending_generation != 0: + var generation := _ack_pending_generation + _ack_pending_generation = 0 + _ack_completed.emit(generation, false) + return true + +func start_auto_flush(owner: Node) -> bool: + if owner == null or not is_instance_valid(owner) or not owner.is_inside_tree(): + return false + stop_auto_flush() + _stopped = false var timer := Timer.new() timer.name = "TelemetryAutoFlushTimer" timer.one_shot = false - timer.autostart = true - timer.wait_time = maxf(_config.batch_interval_s, 0.01) + timer.wait_time = _config.batch_interval_s owner.add_child(timer) timer.timeout.connect(Callable(self, "_on_auto_flush_timeout")) + owner.tree_exiting.connect(Callable(self, "_on_auto_flush_owner_exiting"), CONNECT_ONE_SHOT) + _auto_flush_owner = owner _auto_flush_timer = timer + if _config.enabled: + timer.start() + return true -## Stops timer-based automatic flush. func stop_auto_flush() -> void: + if _auto_flush_owner and is_instance_valid(_auto_flush_owner) and _auto_flush_owner.tree_exiting.is_connected(Callable(self, "_on_auto_flush_owner_exiting")): + _auto_flush_owner.tree_exiting.disconnect(Callable(self, "_on_auto_flush_owner_exiting")) if _auto_flush_timer and is_instance_valid(_auto_flush_timer): _auto_flush_timer.stop() _auto_flush_timer.queue_free() _auto_flush_timer = null + _auto_flush_owner = null -func _on_auto_flush_timeout() -> void: - flush() +func has_auto_flush_timer() -> bool: + return _auto_flush_timer != null and is_instance_valid(_auto_flush_timer) -## Serializes one telemetry event as a dictionary payload. func to_dict(event: TelemetryEvent) -> Dictionary: if event == null: return {} - var serialized_metadata: Dictionary = {} - if event.metadata is Dictionary: - for key in event.metadata.keys(): - serialized_metadata[str(key)] = event.metadata[key] - - return { - "timestamp": event.timestamp_msec, - "level": event.level, - "context_id": event.context_id, - "subject_id": event.subject_id, - "message": event.message, - "metadata": serialized_metadata - } - -## Returns the current telemetry configuration object. + var metadata := {} + for key in event.metadata: + metadata[str(key)] = event.metadata[key] + return {"timestamp": event.timestamp_msec, "level": event.level, "context_id": event.context_id, "subject_id": event.subject_id, "message": event.message, "metadata": metadata} + func get_config() -> TelemetryConfig: return _config + +func _serialize_in_flight() -> Array[Dictionary]: + var output: Array[Dictionary] = [] + for event in _in_flight: + var item := to_dict(event) + if not _is_json_compatible(item): + return [] + output.append(item) + return output + +func _is_json_compatible(value: Variant) -> bool: + match typeof(value): + TYPE_NIL, TYPE_BOOL, TYPE_INT, TYPE_STRING: + return true + TYPE_FLOAT: + return is_finite(value) + TYPE_ARRAY: + for item in value: + if not _is_json_compatible(item): return false + return true + TYPE_DICTIONARY: + for key in value: + if not key is String or not _is_json_compatible(value[key]): return false + return true + _: + return false + +func _restore_in_flight() -> void: + var restored: Array[TelemetryEvent] = [] + restored.append_array(_in_flight) + restored.append_array(_event_queue) + _event_queue = restored + _in_flight.clear() + +func _call_with_timeout(batch: Array[Dictionary]) -> bool: + var response: Variant = _config.flush_callback.call(batch.duplicate(true)) + if response is bool: + return response + if not response is Signal: + return false + _ack_generation += 1 + var generation := _ack_generation + _ack_pending_generation = generation + response.connect(Callable(self, "_on_async_ack").bind(generation), CONNECT_ONE_SHOT) + Engine.get_main_loop().create_timer(_config.callback_timeout_s).timeout.connect(Callable(self, "_on_ack_timeout").bind(generation), CONNECT_ONE_SHOT) + var resolved: Array = await _ack_completed + if resolved[0] != generation: + return false + return resolved[1] + +func _on_async_ack(value: Variant, generation: int) -> void: + if generation == _ack_pending_generation: + _ack_pending_generation = 0 + _ack_completed.emit(generation, value is bool and value) + +func _on_ack_timeout(generation: int) -> void: + if generation == _ack_pending_generation: + _ack_pending_generation = 0 + _counters.callback_timeouts += 1 + _ack_completed.emit(generation, false) + +func _on_auto_flush_timeout() -> void: + if not _flushing: + _flush_requested.emit() + +func _on_auto_flush_owner_exiting() -> void: + if _auto_flush_timer and is_instance_valid(_auto_flush_timer): + _auto_flush_timer.stop() + _auto_flush_timer = null + _auto_flush_owner = null diff --git a/tests/telemetry_module_test.gd b/tests/telemetry_module_test.gd index e5acc7f..7606a21 100644 --- a/tests/telemetry_module_test.gd +++ b/tests/telemetry_module_test.gd @@ -1,55 +1,163 @@ extends SceneTree -const TelemetryModule = preload("res://addon/src/telemetry_module.gd") +const Telemetry = preload("res://addon/src/telemetry_module.gd") +signal callback_ack(value: bool) +signal callback_never(value: bool) -var _flush_call_count: int = 0 -var _last_batch_size: int = 0 +var failures: Array[String] = [] +var callback_mode := "success" +var callback_calls := 0 +var batches: Array = [] func _initialize() -> void: - var failures: Array[String] = [] - await _test_auto_flush_lifecycle(failures) - + await _test_acknowledgement_and_retry() + await _test_queue_pressure_and_pending_isolation() + await _test_duplicate_timeout_serialization_shutdown() + await _test_timer_ownership() if failures.is_empty(): print("TEST_REACHED:telemetry_module_test.gd") print("PASS gd-telemetry telemetry_module_test") quit(0) return - for failure in failures: push_error(failure) quit(1) -func _test_auto_flush_lifecycle(failures: Array[String]) -> void: - _flush_call_count = 0 - _last_batch_size = 0 +func _config(callback: Callable, cap := 1000, batch := 50, retries := 3, timeout := 0.02) -> Telemetry.TelemetryConfig: + return Telemetry.TelemetryConfig.new(true, batch, 0.01, 50, callback, cap, retries, 0.001, 0.002, timeout) - var telemetry := TelemetryModule.new() - var config := TelemetryModule.TelemetryConfig.new(true, 50, 0.01, 50, Callable(self, "_capture_flush")) - telemetry.configure(config) +func _event(id: int, metadata := {}) -> Telemetry.TelemetryEvent: + return Telemetry.TelemetryEvent.new(id, "info", "context", "subject", "event-%d" % id, metadata) - var owner := Node.new() - get_root().add_child(owner) - telemetry.start_auto_flush(owner) +func _check(condition: bool, message: String) -> void: + if not condition: + failures.append(message) - telemetry.add_event(telemetry.build_event(Time.get_ticks_msec(), "info", "session-1", "player-1", "message", {"ok": true})) - if telemetry.event_count() != 1: - failures.append("Expected event_count() to report queued events") - await create_timer(0.05).timeout +func _reset_callback(mode: String) -> void: + callback_mode = mode + callback_calls = 0 + batches.clear() - if _flush_call_count <= 0: - failures.append("Expected auto flush timer to invoke flush callback") - if _last_batch_size != 1: - failures.append("Expected flushed batch to include queued event") +func _callback(batch: Array[Dictionary]) -> Variant: + callback_calls += 1 + batches.append(batch.duplicate(true)) + match callback_mode: + "success": return true + "failure": return false + "invalid": return "not-an-ack" + "async": + call_deferred("_emit_callback_ack", true) + return callback_ack + "pending": return callback_ack + "timeout": return callback_never + return false - telemetry.stop_auto_flush() - telemetry.add_event(telemetry.build_event(Time.get_ticks_msec(), "info", "session-1", "player-1", "message2", {})) - var calls_before_wait: int = _flush_call_count - await create_timer(0.03).timeout - if _flush_call_count != calls_before_wait: - failures.append("Expected stop_auto_flush to stop timer-driven callbacks") +func _emit_callback_ack(value: bool) -> void: + callback_ack.emit(value) - owner.queue_free() +func _test_acknowledgement_and_retry() -> void: + _reset_callback("failure") + var telemetry := Telemetry.new() + _check(telemetry.configure(_config(Callable(self, "_callback"), 10, 2)), "valid configuration rejected") + telemetry.add_event(_event(1)) + var result: Telemetry.FlushResult = await telemetry.flush() + _check(result == Telemetry.FlushResult.FAILED_AFTER_RETRIES, "failed callback did not report exhausted retries") + _check(callback_calls == 4, "expected initial callback plus three retries") + _check(batches.size() == 4 and batches.all(func(batch): return batch[0].timestamp == 1), "retry changed FIFO batch identity") + _check(telemetry.event_count() == 1 and telemetry.in_flight_count() == 0, "failed batch was not restored") + _check(telemetry.counters().retry_attempts == 3, "retry counter mismatch") + _reset_callback("success") + result = await telemetry.flush() + _check(result == Telemetry.FlushResult.ACKNOWLEDGED and telemetry.event_count() == 0, "acknowledged batch was not removed") + +func _test_queue_pressure_and_pending_isolation() -> void: + _reset_callback("invalid") + var telemetry := Telemetry.new() + telemetry.configure(_config(Callable(), 1000, 1000)) + for id in range(1000): telemetry.add_event(_event(id)) + var add_result := telemetry.add_event(_event(1000)) + _check(add_result == Telemetry.AddResult.DROPPED_OLDEST_AND_ADDED, "oldest-first overflow result missing") + _check(telemetry.counters().dropped_oldest == 1 and telemetry.event_count() == 1000, "queue cap/drop counter mismatch") + _check(await telemetry.flush() == Telemetry.FlushResult.INVALID_CALLBACK and telemetry.event_count() == 1000, "invalid callback pressure was not retained/observable") + telemetry.configure(_config(Callable(self, "_callback"), 1000, 1000)) + _reset_callback("success") + await telemetry.flush() + _check(batches[0][0].timestamp == 1 and batches[0][-1].timestamp == 1000, "oldest-first FIFO order mismatch") + + telemetry = Telemetry.new() + telemetry.configure(_config(Callable(), 2, 2)) + telemetry.add_event(_event(10)); telemetry.add_event(_event(11)) + telemetry.configure(_config(Callable(self, "_callback"), 2, 2, 0)) + _reset_callback("pending") + telemetry._flush_requested.emit() + _check(telemetry.in_flight_count() == 2, "batch was not isolated in flight") + _check(telemetry.add_event(_event(12)) == Telemetry.AddResult.REJECTED_IN_FLIGHT_CAPACITY, "newest event was not rejected when only in-flight capacity remained") + _check(telemetry.counters().dropped_newest == 1, "newest rejection counter mismatch") + callback_ack.emit(false) + await process_frame + _check(telemetry.queued_count() == 2 and telemetry.in_flight_count() == 0, "failed pending batch was not restored in isolation") + +func _test_duplicate_timeout_serialization_shutdown() -> void: + var telemetry := Telemetry.new() + telemetry.configure(_config(Callable(self, "_callback"), 3, 2, 0)) + _reset_callback("pending") + telemetry.add_event(_event(1)) + telemetry._flush_requested.emit() + var duplicate: Telemetry.FlushResult = await telemetry.flush() + _check(duplicate == Telemetry.FlushResult.BUSY, "duplicate concurrent flush was not rejected") + callback_ack.emit(true) + await process_frame + _check(telemetry.event_count() == 0 and telemetry.counters().acknowledged_batches == 1, "async true acknowledgement failed") -func _capture_flush(batch: Array[Dictionary]) -> void: - _flush_call_count += 1 - _last_batch_size = batch.size() + _reset_callback("timeout") + telemetry.add_event(_event(2)) + var timed: Telemetry.FlushResult = await telemetry.flush() + _check(timed == Telemetry.FlushResult.FAILED_AFTER_RETRIES and telemetry.counters().callback_timeouts == 1, "callback timeout was not observable") + + telemetry = Telemetry.new() + telemetry.configure(_config(Callable(self, "_callback"), 3, 2, 0)) + _reset_callback("invalid") + telemetry.add_event(_event(20)) + _check(await telemetry.flush() == Telemetry.FlushResult.FAILED_AFTER_RETRIES and telemetry.event_count() == 1, "invalid callback acknowledgement was not a visible failure") + + telemetry = Telemetry.new() + telemetry.configure(_config(Callable(self, "_callback"), 3, 2, 0)) + _reset_callback("success") + telemetry.add_event(_event(3, {"bad": Callable(self, "_callback")})) + _check(await telemetry.flush() == Telemetry.FlushResult.SERIALIZATION_FAILED, "serialization failure was not reported") + _check(telemetry.event_count() == 1 and telemetry.counters().serialization_failures == 1, "serialization failure did not preserve pending event") + _check(telemetry.discard_oldest_event(), "explicit poison-event discard failed") + telemetry.add_event(_event(4)) + _check(await telemetry.shutdown(true) == Telemetry.FlushResult.ACKNOWLEDGED, "shutdown flush outcome mismatch") + _check(telemetry.add_event(_event(5)) == Telemetry.AddResult.DISABLED, "shutdown accepted a new event") + + telemetry = Telemetry.new() + telemetry.configure(_config(Callable(self, "_callback"), 3, 2, 0)) + _reset_callback("pending") + telemetry.add_event(_event(6)) + telemetry._flush_requested.emit() + _check(await telemetry.shutdown(true) == Telemetry.FlushResult.CANCELLED, "shutdown did not cancel and observe an existing flush") + _check(telemetry.event_count() == 1 and telemetry.add_event(_event(7)) == Telemetry.AddResult.DISABLED, "cancelled shutdown lost pending data or stayed enabled") + +func _test_timer_ownership() -> void: + _reset_callback("success") + var telemetry := Telemetry.new() + telemetry.configure(_config(Callable(self, "_callback"), 10, 10)) + var owner := Node.new() + get_root().add_child(owner) + _check(telemetry.start_auto_flush(owner), "valid auto-flush owner rejected") + telemetry.add_event(_event(1)) + await create_timer(0.04).timeout + _check(callback_calls > 0, "auto-flush timer did not invoke callback") + var replacement := _config(Callable(self, "_callback"), 10, 10) + replacement.batch_interval_s = 0.02 + _check(telemetry.configure(replacement) and telemetry.has_auto_flush_timer(), "reconfigure lost timer ownership") + var disabled := _config(Callable(self, "_callback"), 10, 10) + disabled.enabled = false + _check(telemetry.configure(disabled) and telemetry.add_event(_event(2)) == Telemetry.AddResult.DISABLED, "disabled reconfigure still accepted events") + _check(telemetry.configure(replacement) and telemetry.has_auto_flush_timer(), "reenable did not preserve timer ownership") + owner.queue_free() + await process_frame + _check(not telemetry.has_auto_flush_timer(), "owner destruction retained timer") + telemetry.stop_auto_flush() + telemetry.stop_auto_flush()