railiance-telemetry/tests/test_receiver.py

103 lines
5.2 KiB
Python
Raw Normal View History

from datetime import datetime, timezone, timedelta
import json
from pathlib import Path
import sys
import tempfile
import unittest
import uuid
sys.path.insert(0, str(Path(__file__).resolve().parents[1] / 'scripts'))
from receiver import Receiver, strict_json
from platform_event import translate
NOW = datetime(2026, 9, 6, tzinfo=timezone.utc)
class ReceiverTests(unittest.TestCase):
def setUp(self):
self.temp = tempfile.TemporaryDirectory()
self.addCleanup(self.temp.cleanup)
self.path = Path(self.temp.name) / 'receiver.db'
self.contract = dict(schema='railiance-telemetry.stream.v1', stream='test',
producer='platform', recipient='test-operator', signals=['db.restore'],
max_event_age_seconds=900, heartbeat_seconds=900, retention_days=30)
self.receiver = Receiver(self.path, self.contract)
self.addCleanup(lambda: self.receiver.close())
def event(self, state='healthy', when=NOW):
return dict(schema='railiance-telemetry.signal.v1', id=str(uuid.uuid4()),
stream='test', producer='platform', observed_at=when.isoformat(), states={'db.restore': state})
def test_failure_is_durable_and_addressed(self):
self.receiver.ingest(self.event('failed'), NOW)
self.receiver.close()
self.receiver = Receiver(self.path, self.contract)
inbox = self.receiver.inbox()
self.assertEqual(inbox['recipient'], 'test-operator')
self.assertEqual(inbox['notices'][0]['states'], {'db.restore': 'failed'})
self.assertEqual(inbox['notices'][0]['kind'], 'unhealthy')
self.assertTrue(self.receiver.ack(inbox['notices'][0]['id'], NOW)['acknowledged'])
self.assertEqual(self.receiver.inbox()['notices'], [])
def test_never_seen_and_stopped_producer_detected(self):
self.assertEqual(self.receiver.check(NOW)['state'], 'missing-emission')
self.receiver.check(NOW)
self.assertEqual(len(self.receiver.inbox()['notices']), 1)
self.receiver.ingest(self.event(), NOW)
self.assertEqual(self.receiver.check(NOW)['state'], 'receiving')
self.assertEqual(self.receiver.check(NOW + timedelta(seconds=901))['state'], 'missing-emission')
def test_replay_does_not_renew_heartbeat(self):
event = self.event()
self.receiver.ingest(event, NOW)
self.assertEqual(self.receiver.ingest(event, NOW + timedelta(seconds=899))['status'], 'duplicate')
self.assertEqual(self.receiver.check(NOW + timedelta(seconds=901))['state'], 'missing-emission')
with self.assertRaises(ValueError):
self.receiver.ingest(self.event(when=NOW - timedelta(seconds=1)), NOW)
def test_duplicate_id_changed_payload_rejected(self):
event = self.event(); self.receiver.ingest(event, NOW)
event['states']['db.restore'] = 'failed'
with self.assertRaises(ValueError): self.receiver.ingest(event, NOW)
def test_invalid_scope_future_stale_and_extra_data(self):
variants = [dict(producer='other'), dict(states={'other': 'healthy'}),
dict(states={'db.restore': 'secret'}), dict(password='PRIVATE_CANARY'),
dict(observed_at=(NOW + timedelta(seconds=1)).isoformat()),
dict(observed_at=(NOW - timedelta(seconds=901)).isoformat()),
dict(observed_at='2026-09-06T00:00:00')]
for variant in variants:
with self.subTest(variant=variant), self.assertRaises(ValueError):
self.receiver.ingest(dict(self.event(), **variant), NOW)
self.assertEqual(self.receiver.inbox()['notices'], [])
def test_contract_drift_refused(self):
with self.assertRaises(ValueError): Receiver(self.path, dict(self.contract, recipient='other'))
def test_retention_keeps_unacknowledged_and_latest_anchors(self):
self.receiver.ingest(self.event('failed'), NOW)
self.receiver.ingest(self.event(when=NOW + timedelta(seconds=1)), NOW + timedelta(seconds=1))
later = NOW + timedelta(days=31)
self.assertEqual(self.receiver.prune(later)['events_pruned'], 0)
first = self.receiver.inbox()['notices'][0]['id']
self.receiver.ack(first, later)
self.assertEqual(self.receiver.prune(later)['events_pruned'], 1)
self.assertEqual(self.receiver.check(later)['state'], 'missing-emission')
def test_duplicate_fields_and_oversize_rejected(self):
for raw in (b'{"a":1,"a":2}', b' ' * 32769):
with self.assertRaises(ValueError): strict_json(raw)
def test_adapter_keeps_time_and_owner_meaning(self):
report = dict(schema='railiance-platform.assurance-signal.v1',
cluster_uid='a553c742-0115-43d4-99a4-a5ca56fe0786', evaluated_at=NOW.isoformat(),
signals={'db.restore': {'state': 'stale', 'owner': 'railiance-platform'}},
transport='unmonitored', guarantees='unsupported', threshold_status='local-diagnostic-only', healthy=False)
event = translate(report, self.contract)
self.assertEqual(event['observed_at'], report['evaluated_at'])
self.assertEqual(event['states'], {'db.restore': 'stale'})
report['healthy'] = True
with self.assertRaises(ValueError): translate(report, self.contract)
if __name__ == '__main__': unittest.main()