mirror of
https://github.com/dartdavros/chatballs.git
synced 2026-10-06 01:24:59 +03:00
feat(calls): call-метрики без медиаконтента — проход D этапа E09
- CallMetric: тип соединения DIRECT/RELAY/UNKNOWN + категория ICE-кандидата (host/srflx/prflx/relay) и RTT; храним только КАТЕГОРИЮ, без адресов/SDP/ICE - record_call_metric: derive connection_type, whitelist кандидатов, clamp RTT, идемпотентный upsert по (call_session, side) - signaling/consumer: тип participant.metrics — сохраняем, не ретранслируем и не логируем; call_payload отдаёт metrics для internal-ui/аналитики E15 - callRtc.ts: на connected снимает getStats выбранной candidate-pair и шлёт только категорию кандидата + RTT (подтверждает direct vs TURN relay) - tests: derive/relay/unknown, санитизация мусора и RTT, upsert, payload Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
1 parent
e8183777df
commit
c55a3199fe
9 files changed
+288
-1
No files matched your search
@@ -1,6 +1,6 @@
|
||||
from django.contrib import admin
|
||||
|
||||
from hub_platform.calls.models import CallInvite, CallParticipant, CallSession
|
||||
from hub_platform.calls.models import CallInvite, CallMetric, CallParticipant, CallSession
|
||||
|
||||
|
||||
@admin.register(CallSession)
|
||||
@@ -21,3 +21,9 @@ class CallInviteAdmin(admin.ModelAdmin):
|
||||
class CallParticipantAdmin(admin.ModelAdmin):
|
||||
list_display = ["call_session", "side", "last_connection_state", "joined_at", "left_at"]
|
||||
list_filter = ["side", "last_connection_state"]
|
||||
|
||||
|
||||
@admin.register(CallMetric)
|
||||
class CallMetricAdmin(admin.ModelAdmin):
|
||||
list_display = ["call_session", "side", "connection_type", "round_trip_ms", "updated_at"]
|
||||
list_filter = ["connection_type", "side"]
|
||||
@@ -54,6 +54,9 @@ class CallSignalingConsumer(AsyncJsonWebsocketConsumer):
|
||||
await self._relay(msg_type, content)
|
||||
elif msg_type == "participant.connection_state":
|
||||
await self._connection_state(content)
|
||||
elif msg_type == "participant.metrics":
|
||||
# Технические метрики без медиаконтента: сохраняем, не ретранслируем.
|
||||
await database_sync_to_async(signaling.record_metric)(self.call_id, self.side, content)
|
||||
elif msg_type == "call.ended":
|
||||
payload = await database_sync_to_async(signaling.end_from_signaling)(self.call_id, self.side)
|
||||
await self._broadcast({"type": "call.state", "call": payload}, include_self=True)
|
||||
|
||||
@@ -0,0 +1,31 @@
|
||||
# Generated by Django 5.2.16 on 2026-07-12 01:32
|
||||
|
||||
import django.db.models.deletion
|
||||
from django.db import migrations, models
|
||||
|
||||
|
||||
class Migration(migrations.Migration):
|
||||
|
||||
dependencies = [
|
||||
('calls', '0001_initial'),
|
||||
]
|
||||
|
||||
operations = [
|
||||
migrations.CreateModel(
|
||||
name='CallMetric',
|
||||
fields=[
|
||||
('id', models.BigAutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')),
|
||||
('side', models.CharField(choices=[('STAFF', 'Сотрудник'), ('CUSTOMER', 'Клиент')], max_length=16)),
|
||||
('connection_type', models.CharField(choices=[('DIRECT', 'Прямое P2P'), ('RELAY', 'Через TURN relay'), ('UNKNOWN', 'Неизвестно')], default='UNKNOWN', max_length=16)),
|
||||
('local_candidate_type', models.CharField(blank=True, max_length=8)),
|
||||
('remote_candidate_type', models.CharField(blank=True, max_length=8)),
|
||||
('round_trip_ms', models.PositiveIntegerField(blank=True, null=True)),
|
||||
('created_at', models.DateTimeField(auto_now_add=True)),
|
||||
('updated_at', models.DateTimeField(auto_now=True)),
|
||||
('call_session', models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name='metrics', to='calls.callsession')),
|
||||
],
|
||||
options={
|
||||
'constraints': [models.UniqueConstraint(fields=('call_session', 'side'), name='uniq_call_metric_side')],
|
||||
},
|
||||
),
|
||||
]
|
||||
@@ -62,6 +62,12 @@ class InviteDeliveryStatus(models.TextChoices):
|
||||
FAILED = "FAILED", "Ошибка доставки"
|
||||
|
||||
|
||||
class CallConnectionType(models.TextChoices):
|
||||
DIRECT = "DIRECT", "Прямое P2P"
|
||||
RELAY = "RELAY", "Через TURN relay"
|
||||
UNKNOWN = "UNKNOWN", "Неизвестно"
|
||||
|
||||
|
||||
class CallSession(models.Model):
|
||||
id = models.UUIDField(primary_key=True, default=uuid.uuid4, editable=False)
|
||||
organization = models.ForeignKey("identity.Organization", on_delete=models.PROTECT, related_name="call_sessions")
|
||||
@@ -171,3 +177,34 @@ class CallParticipant(models.Model):
|
||||
|
||||
def __str__(self) -> str:
|
||||
return f"participant:{self.call_session_id}/{self.side}"
|
||||
|
||||
|
||||
class CallMetric(models.Model):
|
||||
"""Технические метрики соединения без медиаконтента (SPEC-HUB-0013 §13).
|
||||
|
||||
Хранится только КАТЕГОРИЯ ICE-кандидата (host/srflx/prflx/relay) и RTT, но
|
||||
никогда сам ICE candidate, его адрес, SDP или медиапоток. Позволяет считать
|
||||
долю direct/relay звонков (E15) и подтверждать TURN fallback, не раскрывая
|
||||
сетевые адреса участников.
|
||||
"""
|
||||
|
||||
call_session = models.ForeignKey(CallSession, on_delete=models.CASCADE, related_name="metrics")
|
||||
side = models.CharField(max_length=16, choices=ParticipantSide.choices)
|
||||
connection_type = models.CharField(
|
||||
max_length=16,
|
||||
choices=CallConnectionType.choices,
|
||||
default=CallConnectionType.UNKNOWN,
|
||||
)
|
||||
local_candidate_type = models.CharField(max_length=8, blank=True)
|
||||
remote_candidate_type = models.CharField(max_length=8, blank=True)
|
||||
round_trip_ms = models.PositiveIntegerField(null=True, blank=True)
|
||||
created_at = models.DateTimeField(auto_now_add=True)
|
||||
updated_at = models.DateTimeField(auto_now=True)
|
||||
|
||||
class Meta:
|
||||
constraints = [
|
||||
models.UniqueConstraint(fields=["call_session", "side"], name="uniq_call_metric_side"),
|
||||
]
|
||||
|
||||
def __str__(self) -> str:
|
||||
return f"metric:{self.call_session_id}/{self.side}/{self.connection_type}"
|
||||
@@ -18,6 +18,16 @@ def call_payload(call: CallSession) -> dict:
|
||||
}
|
||||
for participant in call.participants.all()
|
||||
]
|
||||
metrics = [
|
||||
{
|
||||
"side": metric.side,
|
||||
"connectionType": metric.connection_type,
|
||||
"localCandidateType": metric.local_candidate_type or None,
|
||||
"remoteCandidateType": metric.remote_candidate_type or None,
|
||||
"roundTripMs": metric.round_trip_ms,
|
||||
}
|
||||
for metric in call.metrics.all()
|
||||
]
|
||||
return {
|
||||
"id": str(call.id),
|
||||
"conversationId": call.conversation_id,
|
||||
@@ -32,6 +42,7 @@ def call_payload(call: CallSession) -> dict:
|
||||
"failureCode": call.failure_code or None,
|
||||
"durationSeconds": call.duration_seconds,
|
||||
"participants": participants,
|
||||
"metrics": metrics,
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -10,8 +10,10 @@ from django.utils import timezone
|
||||
from hub_platform.calls.errors import CallConflict, CallInvalidTransition, CallTokenError
|
||||
from hub_platform.calls.lifecycle import transition_call
|
||||
from hub_platform.calls.models import (
|
||||
CallConnectionType,
|
||||
CallEndedBy,
|
||||
CallInvite,
|
||||
CallMetric,
|
||||
CallParticipant,
|
||||
CallSession,
|
||||
CallStatus,
|
||||
@@ -287,6 +289,60 @@ def call_state_by_access_token(*, token: str) -> CallSession:
|
||||
return call
|
||||
|
||||
|
||||
# ICE candidate типы (RFC 8445), допустимые в метриках. Храним только категорию,
|
||||
# без адреса/порта/foundation самого кандидата.
|
||||
_ALLOWED_CANDIDATE_TYPES = {"host", "srflx", "prflx", "relay"}
|
||||
_MAX_ROUND_TRIP_MS = 60_000
|
||||
|
||||
|
||||
def _sanitize_candidate_type(value) -> str:
|
||||
text = str(value or "").lower()
|
||||
return text if text in _ALLOWED_CANDIDATE_TYPES else ""
|
||||
|
||||
|
||||
def _connection_type(local: str, remote: str) -> str:
|
||||
if "relay" in (local, remote):
|
||||
return CallConnectionType.RELAY
|
||||
if local or remote:
|
||||
return CallConnectionType.DIRECT
|
||||
return CallConnectionType.UNKNOWN
|
||||
|
||||
|
||||
def record_call_metric(
|
||||
*,
|
||||
call_session_id,
|
||||
side: str,
|
||||
local_candidate_type=None,
|
||||
remote_candidate_type=None,
|
||||
round_trip_ms=None,
|
||||
) -> None:
|
||||
"""Сохранить технические метрики соединения участника (без медиаконтента).
|
||||
|
||||
Идемпотентно по (call_session, side): при reconnect/ICE-restart метрика
|
||||
обновляется актуальным типом маршрута. Никакие SDP/ICE payload не пишутся —
|
||||
только производная категория кандидата и RTT.
|
||||
"""
|
||||
if side not in ParticipantSide.values:
|
||||
return
|
||||
if not CallSession.objects.filter(id=call_session_id).exists():
|
||||
return
|
||||
local = _sanitize_candidate_type(local_candidate_type)
|
||||
remote = _sanitize_candidate_type(remote_candidate_type)
|
||||
rtt: int | None = None
|
||||
if isinstance(round_trip_ms, (int, float)) and not isinstance(round_trip_ms, bool):
|
||||
rtt = max(0, min(_MAX_ROUND_TRIP_MS, int(round_trip_ms)))
|
||||
CallMetric.objects.update_or_create(
|
||||
call_session_id=call_session_id,
|
||||
side=side,
|
||||
defaults={
|
||||
"connection_type": _connection_type(local, remote),
|
||||
"local_candidate_type": local,
|
||||
"remote_candidate_type": remote,
|
||||
"round_trip_ms": rtt,
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
def active_call_for_conversation(conversation: Conversation) -> CallSession | None:
|
||||
return (
|
||||
CallSession.objects.filter(
|
||||
|
||||
@@ -19,6 +19,7 @@ from hub_platform.calls.models import (
|
||||
TERMINAL_CALL_STATUSES,
|
||||
)
|
||||
from hub_platform.calls.serializers import public_call_state_payload
|
||||
from hub_platform.calls.services import record_call_metric
|
||||
|
||||
SIDE_TO_ENDED_BY = {
|
||||
ParticipantSide.STAFF: CallEndedBy.STAFF,
|
||||
@@ -123,6 +124,21 @@ def report_connection(call_id, side: str, connected: bool) -> tuple[dict | None,
|
||||
return public_call_state_payload(call), True
|
||||
|
||||
|
||||
def record_metric(call_id, side: str, content: dict) -> None:
|
||||
"""Метрики соединения от участника: только категория маршрута и RTT.
|
||||
|
||||
Ничего не ретранслируется собеседнику и не логируется — payload не содержит
|
||||
медиаконтента, но и типы кандидатов наружу не пересылаются.
|
||||
"""
|
||||
record_call_metric(
|
||||
call_session_id=call_id,
|
||||
side=side,
|
||||
local_candidate_type=content.get("localCandidateType"),
|
||||
remote_candidate_type=content.get("remoteCandidateType"),
|
||||
round_trip_ms=content.get("roundTripMs"),
|
||||
)
|
||||
|
||||
|
||||
def end_from_signaling(call_id, side: str) -> dict:
|
||||
"""Завершение звонка стороной: идемпотентно, целевой статус — по фазе."""
|
||||
ended_by = SIDE_TO_ENDED_BY.get(side, CallEndedBy.SYSTEM)
|
||||
|
||||
@@ -0,0 +1,81 @@
|
||||
import uuid
|
||||
|
||||
from hub_platform.calls import signaling
|
||||
from hub_platform.calls.models import CallConnectionType, CallMetric, ParticipantSide
|
||||
from hub_platform.calls.serializers import call_payload
|
||||
from hub_platform.calls.services import create_call_request, record_call_metric
|
||||
from hub_platform.calls.tests.helpers import CallTestCase
|
||||
|
||||
|
||||
class RecordCallMetricTests(CallTestCase):
|
||||
def setUp(self) -> None:
|
||||
super().setUp()
|
||||
self.call = create_call_request(
|
||||
conversation_id=self.conversation.id, initiator=self.owner
|
||||
).call_session
|
||||
|
||||
def _record(self, **kwargs):
|
||||
record_call_metric(call_session_id=self.call.id, side=ParticipantSide.STAFF, **kwargs)
|
||||
return CallMetric.objects.get(call_session=self.call, side=ParticipantSide.STAFF)
|
||||
|
||||
def test_direct_connection_type_from_non_relay_candidates(self) -> None:
|
||||
metric = self._record(local_candidate_type="host", remote_candidate_type="srflx", round_trip_ms=42)
|
||||
self.assertEqual(metric.connection_type, CallConnectionType.DIRECT)
|
||||
self.assertEqual(metric.local_candidate_type, "host")
|
||||
self.assertEqual(metric.round_trip_ms, 42)
|
||||
|
||||
def test_relay_candidate_marks_turn_fallback(self) -> None:
|
||||
metric = self._record(local_candidate_type="relay", remote_candidate_type="srflx")
|
||||
self.assertEqual(metric.connection_type, CallConnectionType.RELAY)
|
||||
|
||||
def test_unknown_when_no_candidate_types(self) -> None:
|
||||
metric = self._record(local_candidate_type=None, remote_candidate_type=None)
|
||||
self.assertEqual(metric.connection_type, CallConnectionType.UNKNOWN)
|
||||
|
||||
def test_invalid_candidate_type_is_dropped(self) -> None:
|
||||
# Клиент прислал мусор/адрес — храним только валидную категорию (RFC 8445).
|
||||
metric = self._record(local_candidate_type="203.0.113.7", remote_candidate_type="RELAY")
|
||||
self.assertEqual(metric.local_candidate_type, "")
|
||||
self.assertEqual(metric.remote_candidate_type, "relay")
|
||||
self.assertEqual(metric.connection_type, CallConnectionType.RELAY)
|
||||
|
||||
def test_round_trip_ms_sanitized(self) -> None:
|
||||
self.assertEqual(self._record(round_trip_ms=-5).round_trip_ms, 0)
|
||||
self.assertEqual(self._record(round_trip_ms=10**9).round_trip_ms, 60_000)
|
||||
self.assertEqual(self._record(round_trip_ms=12.9).round_trip_ms, 12)
|
||||
self.assertIsNone(self._record(round_trip_ms=True).round_trip_ms)
|
||||
self.assertIsNone(self._record(round_trip_ms="fast").round_trip_ms)
|
||||
|
||||
def test_upsert_is_idempotent_per_side(self) -> None:
|
||||
self._record(local_candidate_type="host")
|
||||
self._record(local_candidate_type="relay")
|
||||
metrics = CallMetric.objects.filter(call_session=self.call, side=ParticipantSide.STAFF)
|
||||
self.assertEqual(metrics.count(), 1)
|
||||
self.assertEqual(metrics.first().connection_type, CallConnectionType.RELAY)
|
||||
|
||||
def test_unknown_side_ignored(self) -> None:
|
||||
record_call_metric(call_session_id=self.call.id, side="ALIEN", local_candidate_type="host")
|
||||
self.assertFalse(CallMetric.objects.filter(call_session=self.call).exists())
|
||||
|
||||
def test_missing_call_ignored(self) -> None:
|
||||
record_call_metric(call_session_id=uuid.uuid4(), side=ParticipantSide.STAFF, local_candidate_type="host")
|
||||
self.assertEqual(CallMetric.objects.count(), 0)
|
||||
|
||||
def test_signaling_helper_maps_camel_case_payload(self) -> None:
|
||||
signaling.record_metric(
|
||||
self.call.id,
|
||||
ParticipantSide.CUSTOMER,
|
||||
{"localCandidateType": "relay", "remoteCandidateType": "host", "roundTripMs": 88},
|
||||
)
|
||||
metric = CallMetric.objects.get(call_session=self.call, side=ParticipantSide.CUSTOMER)
|
||||
self.assertEqual(metric.connection_type, CallConnectionType.RELAY)
|
||||
self.assertEqual(metric.round_trip_ms, 88)
|
||||
|
||||
def test_call_payload_exposes_metrics(self) -> None:
|
||||
self._record(local_candidate_type="host", remote_candidate_type="host", round_trip_ms=15)
|
||||
payload = call_payload(self.call)
|
||||
self.assertEqual(len(payload["metrics"]), 1)
|
||||
entry = payload["metrics"][0]
|
||||
self.assertEqual(entry["side"], ParticipantSide.STAFF)
|
||||
self.assertEqual(entry["connectionType"], CallConnectionType.DIRECT)
|
||||
self.assertEqual(entry["roundTripMs"], 15)
|
||||
@@ -192,6 +192,7 @@ export class CallRtcClient {
|
||||
case "connected":
|
||||
this.sendCommand({ type: "participant.connection_state", state: "CONNECTED" });
|
||||
this.options.handlers.onConnection?.("connected");
|
||||
void this.reportMetrics();
|
||||
break;
|
||||
case "disconnected":
|
||||
this.sendCommand({ type: "participant.connection_state", state: "RECONNECTING" });
|
||||
@@ -211,6 +212,51 @@ export class CallRtcClient {
|
||||
return pc;
|
||||
}
|
||||
|
||||
/**
|
||||
* Технические метрики соединения (SPEC §13): только КАТЕГОРИЯ выбранного
|
||||
* ICE-кандидата (host/srflx/relay) и RTT — чтобы отличить direct от TURN relay.
|
||||
* Ни SDP, ни адреса кандидатов, ни медиаданные не отправляются.
|
||||
*/
|
||||
private async reportMetrics(): Promise<void> {
|
||||
const pc = this.pc;
|
||||
if (!pc) return;
|
||||
try {
|
||||
const stats = await pc.getStats();
|
||||
const get = (id: unknown): Record<string, unknown> | undefined =>
|
||||
typeof id === "string" ? (stats.get(id) as Record<string, unknown> | undefined) : undefined;
|
||||
let selectedId = "";
|
||||
let pair: Record<string, unknown> | null = null;
|
||||
stats.forEach((report) => {
|
||||
const entry = report as Record<string, unknown>;
|
||||
if (entry.type === "transport" && typeof entry.selectedCandidatePairId === "string") {
|
||||
selectedId = entry.selectedCandidatePairId;
|
||||
}
|
||||
});
|
||||
stats.forEach((report) => {
|
||||
const entry = report as Record<string, unknown>;
|
||||
if (entry.type !== "candidate-pair") return;
|
||||
if (entry.id === selectedId || (!pair && entry.nominated === true && entry.state === "succeeded")) {
|
||||
pair = entry;
|
||||
}
|
||||
});
|
||||
const selected = pair as Record<string, unknown> | null;
|
||||
if (!selected) return;
|
||||
const candidateType = (id: unknown): string => {
|
||||
const type = get(id)?.candidateType;
|
||||
return typeof type === "string" ? type : "";
|
||||
};
|
||||
const rtt = selected.currentRoundTripTime;
|
||||
this.sendCommand({
|
||||
type: "participant.metrics",
|
||||
localCandidateType: candidateType(selected.localCandidateId),
|
||||
remoteCandidateType: candidateType(selected.remoteCandidateId),
|
||||
roundTripMs: typeof rtt === "number" ? Math.round(rtt * 1000) : null,
|
||||
});
|
||||
} catch {
|
||||
/* getStats недоступен/прерван — метрики необязательны */
|
||||
}
|
||||
}
|
||||
|
||||
private publishMediaState(): void {
|
||||
const audio = this.options.localStream?.getAudioTracks()[0];
|
||||
const video = this.options.localStream?.getVideoTracks()[0];
|
||||
|
||||
Reference in new issue
Block a user