TelemetryReader(db: EntityDatabase, app_name: str, num_processes: int, rabbit_config: dict)
Reader of telemetry data.
Used by API.
Not contained inside Telemetry class due to usage of CallbackRegistrar
and all of its requirements.
Source code in dp3/history_management/telemetry.py
| def __init__(
self,
db: EntityDatabase,
app_name: str,
num_processes: int,
rabbit_config: dict,
) -> None:
self.db = db
self.app_name = app_name
self.num_processes = num_processes
self.rabbit_config = rabbit_config or {}
self.cache_col = self.db.get_module_cache("Telemetry")
|
get_sources_validity
get_sources_validity() -> dict[str, datetime]
Return timestamps (datetimes) of current validity of all sources.
Source code in dp3/history_management/telemetry.py
| def get_sources_validity(self) -> dict[str, datetime]:
"""Return timestamps (datetimes) of current validity of all sources."""
src_data = self.cache_col.find({"src_t": {"$exists": True}}).sort([("_id", ASCENDING)])
return {src["_id"]: src["src_t"] for src in src_data}
|
get_sources_age
get_sources_age(unit: str = 'minutes') -> dict[str, int]
Return ages of sources relative to now in requested units.
Source code in dp3/history_management/telemetry.py
| def get_sources_age(self, unit: str = "minutes") -> dict[str, int]:
"""Return ages of sources relative to now in requested units."""
now = datetime.now(UTC)
divider = 60 if unit == "minutes" else 1
return {
source: int((now - timestamp).total_seconds() / divider)
for source, timestamp in self.get_sources_validity().items()
}
|
get_entities_per_attr
get_entities_per_attr() -> dict[str, dict[str, int]]
Return counts of entities with data present for each configured attribute.
Source code in dp3/history_management/telemetry.py
| def get_entities_per_attr(self) -> dict[str, dict[str, int]]:
"""Return counts of entities with data present for each configured attribute."""
return self.db.count_entities_per_attr()
|
get_attribute_bson_sizes
get_attribute_bson_sizes() -> dict[str, dict]
Return the latest cached attribute BSON-size statistics.
Source code in dp3/history_management/telemetry.py
| def get_attribute_bson_sizes(self) -> dict[str, dict]:
"""Return the latest cached attribute BSON-size statistics."""
records = self.cache_col.find({"telemetry_type": ATTRIBUTE_BSON_SIZES_TYPE}).sort(
[("entity_type", ASCENDING)]
)
return {
record["entity_type"]: {
"calculated_at": record["calculated_at"],
"duration_s": record["duration_s"],
"attributes": record["attributes"],
}
for record in records
}
|
get_snapshot_summary
get_snapshot_summary() -> dict
Return summary of latest snapshot activity.
Source code in dp3/history_management/telemetry.py
| def get_snapshot_summary(self) -> dict:
"""Return summary of latest snapshot activity."""
now = datetime.now(UTC)
latest_started = next(self.db.find_metadata(module="SnapShooter").limit(1), None)
latest_finished = next(
self.db.find_metadata(
module="SnapShooter",
extra_filter={"workers_finished": self.num_processes, "linked_finished": True},
).limit(1),
None,
)
summary = {
"latest_age": None,
"finished_age": None,
"entities": None,
"total_s": None,
}
if latest_started is not None:
summary["latest_age"] = (now - latest_started["#time_created"]).total_seconds()
if latest_finished is not None:
summary["finished_age"] = (now - latest_finished["#time_created"]).total_seconds()
summary["entities"] = latest_finished.get("entities")
summary["total_s"] = (
latest_finished["#last_update"] - latest_finished["#time_created"]
).total_seconds()
return summary
|
get_metadata
get_metadata(module: str = None, date_from: datetime = None, date_to: datetime = None, newest_first: bool = True, skip: int = 0, limit: int = 0) -> list[dict]
Return filtered metadata documents.
Source code in dp3/history_management/telemetry.py
| def get_metadata(
self,
module: str = None,
date_from: datetime = None,
date_to: datetime = None,
newest_first: bool = True,
skip: int = 0,
limit: int = 0,
) -> list[dict]:
"""Return filtered metadata documents."""
cursor = self.db.find_metadata(module, date_from, date_to, newest_first)
if skip:
cursor = cursor.skip(skip)
if limit:
cursor = cursor.limit(limit)
return list(cursor)
|
get_rabbitmq_queues
get_rabbitmq_queues() -> dict[str, list[dict]]
Return RabbitMQ queue telemetry for this application.
Source code in dp3/history_management/telemetry.py
| def get_rabbitmq_queues(self) -> dict[str, list[dict]]:
"""Return RabbitMQ queue telemetry for this application."""
host = self.rabbit_config.get("host", "localhost")
port = int(self.rabbit_config.get("management_port", 15672))
username = self.rabbit_config.get("username", "guest")
password = self.rabbit_config.get("password", "guest")
response = requests.get(
f"http://{host}:{port}/api/queues",
auth=(username, password),
timeout=5,
)
response.raise_for_status()
queues = []
app_prefix = f"{self.app_name}-worker-"
for queue in response.json():
queue_name = queue.get("name", "")
if not queue_name.startswith(app_prefix):
continue
short_name = queue_name[len(app_prefix) :]
alias = short_name if len(short_name) > 2 else f"{short_name}-main"
queue_data = {"name": queue_name, "queue": alias}
for key, alias_key in self.exported_queue_keys.items():
queue_data[alias_key] = queue.get(key, 0)
message_stats = queue.get("message_stats", {})
for key, alias_key in self.exported_message_stats_keys.items():
queue_data[alias_key] = message_stats.get(key, {}).get("rate", 0)
queues.append(queue_data)
queues.sort(key=lambda item: item["queue"])
return {"queues": queues}
|