Files
truf-server/tests/test_operations_service.py
2026-09-30 20:30:56 +03:00

1053 lines
44 KiB
Python

import asyncio
import json
import os
from pathlib import Path
import sys
import tempfile
import threading
import unittest
import uuid
from unittest import mock
ROOT = Path(__file__).resolve().parents[1]
APP_DIR = ROOT / 'app'
sys.path.insert(0, str(APP_DIR))
import scanner_db
from scanner_db import ScannerDB
class OperationsServiceTests(unittest.TestCase):
def setUp(self):
self.environment = mock.patch.dict(
os.environ, {'SCANNER_DB_URL': '', 'DATABASE_URL': ''},
)
self.environment.start()
self.temp = tempfile.TemporaryDirectory()
self.path = os.path.join(self.temp.name, 'scanner.db')
self.db = ScannerDB(db_path=self.path)
def tearDown(self):
self.db.close()
self.temp.cleanup()
self.environment.stop()
@staticmethod
def identity(action='apply-both'):
return {
'active_config_sha256': '1' * 64,
'active_secrets_sha256': '2' * 64,
'candidate_config_sha256': (
'3' * 64 if action in ('apply-config', 'apply-both') else None
),
'candidate_secrets_sha256': (
'4' * 64 if action in ('apply-secrets', 'apply-both') else None
),
}
@staticmethod
def envelope(operation_id, action, result, *, identity=None, category=None, detail=None):
return json.dumps({
'schema': 1,
'operation_id': operation_id,
'action': action,
'result': result,
'safe_category': category,
'safe_detail': detail,
'resulting_identity': identity,
}, sort_keys=True, separators=(',', ':')).encode('utf-8')
def create(self, action='apply-both'):
operation_id = str(uuid.uuid4())
return operation_id, self.db.create_runtime_operation(
operation_id=operation_id,
actor='operator:alice',
action=action,
expected_identity=self.identity(action),
)
def test_create_replay_commits_requested_operation_and_accepted_audit(self):
operation_id = str(uuid.uuid4())
expected = self.identity()
created = self.db.create_runtime_operation(
operation_id=operation_id, actor='operator:alice',
action='apply-both', expected_identity=expected,
)
replay = self.db.create_runtime_operation(
operation_id=operation_id, actor='operator:alice',
action='apply-both', expected_identity=dict(expected),
)
self.assertEqual(created['status'], 'requested')
self.assertEqual(created['agent_state'], 'pending')
self.assertFalse(created['replayed'])
self.assertTrue(replay['replayed'])
self.assertEqual(replay['audit_event_id'], created['audit_event_id'])
self.assertEqual(replay['expected_identity'], expected)
counts = self.db.conn.execute(
'''SELECT
(SELECT COUNT(*) FROM runtime_operations) AS operations,
(SELECT COUNT(*) FROM runtime_audit_events) AS audits'''
).fetchone()
self.assertEqual((counts['operations'], counts['audits']), (1, 1))
audit = self.db.conn.execute(
'SELECT * FROM runtime_audit_events WHERE operation_id = ?',
(operation_id,),
).fetchone()
self.assertEqual(audit['result'], 'accepted')
self.assertEqual(audit['action'], 'apply-both')
self.assertEqual(audit['target_kind'], 'runtime-deployment')
self.assertEqual(audit['target_ref'], 'config-secrets')
self.assertEqual(json.loads(audit['before_identity_json']), expected)
self.assertIsNone(audit['after_identity_json'])
def test_running_success_and_exact_result_replay_survive_reopen(self):
operation_id, _ = self.create('apply-config')
running = self.db.mark_runtime_operation_running(operation_id)
running_replay = self.db.mark_runtime_operation_running(operation_id)
resulting = {
'active_config_sha256': '3' * 64,
'active_secrets_sha256': '2' * 64,
}
envelope = self.envelope(
operation_id, 'apply-config', 'succeeded', identity=resulting,
)
completed = self.db.reconcile_runtime_operation_result(envelope)
replay = self.db.reconcile_runtime_operation_result(envelope)
self.assertEqual(running['status'], 'running')
self.assertFalse(running['replayed'])
self.assertTrue(running_replay['replayed'])
self.assertEqual(completed['status'], 'succeeded')
self.assertEqual(completed['agent_state'], 'succeeded')
self.assertEqual(completed['resulting_identity'], resulting)
self.assertIsNotNone(completed['started_at'])
self.assertIsNotNone(completed['completed_at'])
self.assertIsNotNone(completed['agent_reconciled_at'])
self.assertTrue(replay['replayed'])
self.assertEqual(replay['agent_result_sha256'], completed['agent_result_sha256'])
self.assertEqual(replay['audit_event_id'], completed['audit_event_id'])
events = self.db.conn.execute(
'''SELECT result, previous_event_id, previous_event_sha256, event_sha256
FROM runtime_audit_events WHERE operation_id = ? ORDER BY id''',
(operation_id,),
).fetchall()
self.assertEqual([row['result'] for row in events], ['accepted', 'succeeded'])
self.assertIsNotNone(events[1]['previous_event_id'])
self.assertEqual(events[1]['previous_event_sha256'], events[0]['event_sha256'])
self.db.close()
self.db = ScannerDB(db_path=self.path)
persisted = self.db.runtime_operation(operation_id)
self.assertEqual(persisted['status'], 'succeeded')
self.assertEqual(persisted['resulting_identity'], resulting)
def test_host_claim_binds_persisted_identity_and_exactly_replays(self):
operation_id, _created = self.create('apply-both')
expected = self.identity('apply-both')
claimed = self.db.claim_runtime_operation_execution(
operation_id=operation_id,
action='apply-both',
expected_identity=expected,
)
replay = self.db.claim_runtime_operation_execution(
operation_id=operation_id,
action='apply-both',
expected_identity=dict(expected),
)
self.assertEqual(claimed['status'], 'running')
self.assertEqual(claimed['agent_state'], 'running')
self.assertFalse(claimed['replayed'])
self.assertTrue(replay['replayed'])
self.assertEqual(replay['audit_event_id'], claimed['audit_event_id'])
self.assertEqual(
replay['audit_event_sha256'], claimed['audit_event_sha256'],
)
events = self.db.conn.execute(
'SELECT result FROM runtime_audit_events WHERE operation_id = ?',
(operation_id,),
).fetchall()
self.assertEqual([row['result'] for row in events], ['accepted'])
def test_host_claim_rejects_request_identity_drift_without_transition(self):
operation_id, _created = self.create('apply-config')
wrong_hash = self.identity('apply-config')
wrong_hash['active_config_sha256'] = 'f' * 64
with self.assertRaises(scanner_db.RuntimeOperationIdentityConflictError):
self.db.claim_runtime_operation_execution(
operation_id=operation_id,
action='apply-config',
expected_identity=wrong_hash,
)
with self.assertRaises(scanner_db.RuntimeOperationIdentityConflictError):
self.db.claim_runtime_operation_execution(
operation_id=operation_id,
action='apply-secrets',
expected_identity=self.identity('apply-secrets'),
)
operation = self.db.runtime_operation(operation_id)
self.assertEqual(operation['status'], 'requested')
self.assertEqual(operation['agent_state'], 'pending')
def test_host_claim_requires_valid_accepted_audit_evidence(self):
operation_id, _created = self.create('apply-config')
self.db.conn.execute('DROP TRIGGER runtime_audit_events_reject_update')
self.db.conn.execute(
'''UPDATE runtime_audit_events SET before_identity_json = ?
WHERE operation_id = ?''',
('{}', operation_id),
)
self.db.conn.commit()
with self.assertRaisesRegex(
scanner_db.RuntimeSafetySchemaError, 'accepted audit evidence',
):
self.db.claim_runtime_operation_execution(
operation_id=operation_id,
action='apply-config',
expected_identity=self.identity('apply-config'),
)
operation = self.db.runtime_operation(operation_id)
self.assertEqual(operation['status'], 'requested')
def test_host_claim_rejects_corrupt_immediate_audit_predecessor(self):
predecessor_id, _created = self.create('restart')
operation_id, _created = self.create('apply-config')
self.db.conn.execute('DROP TRIGGER runtime_audit_events_reject_update')
self.db.conn.execute(
'''UPDATE runtime_audit_events SET actor = ?
WHERE operation_id = ?''',
('actor:tampered', predecessor_id),
)
self.db.conn.commit()
with self.assertRaisesRegex(
scanner_db.RuntimeSafetySchemaError, 'accepted audit evidence',
):
self.db.claim_runtime_operation_execution(
operation_id=operation_id,
action='apply-config',
expected_identity=self.identity('apply-config'),
)
operation = self.db.runtime_operation(operation_id)
self.assertEqual(operation['status'], 'requested')
def test_host_claim_accepts_valid_unlinked_audit_predecessor(self):
self.db._append_runtime_audit_event_locked(
operation_id=None,
actor='system:storage',
action='storage.inspect',
target_kind='storage',
target_ref='',
result='accepted',
)
self.db.conn.commit()
operation_id, _created = self.create('apply-config')
claimed = self.db.claim_runtime_operation_execution(
operation_id=operation_id,
action='apply-config',
expected_identity=self.identity('apply-config'),
)
self.assertEqual(claimed['status'], 'running')
self.assertFalse(claimed['replayed'])
def test_host_claim_rejects_missing_or_terminal_operation(self):
with self.assertRaises(scanner_db.RuntimeOperationTransitionError):
self.db.claim_runtime_operation_execution(
operation_id=str(uuid.uuid4()),
action='restart',
expected_identity=self.identity('restart'),
)
operation_id, _created = self.create('restart')
self.db.mark_runtime_operation_running(operation_id)
self.db.reconcile_runtime_operation_result(self.envelope(
operation_id,
'restart',
'succeeded',
identity={
'active_config_sha256': '1' * 64,
'active_secrets_sha256': '2' * 64,
},
))
with self.assertRaises(scanner_db.RuntimeOperationTransitionError):
self.db.claim_runtime_operation_execution(
operation_id=operation_id,
action='restart',
expected_identity=self.identity('restart'),
)
def test_recent_operations_are_bounded_and_deterministic(self):
first_id, _ = self.create('restart')
second_id, _ = self.create('apply-config')
self.db.conn.execute(
'UPDATE runtime_operations SET updated_at = ? WHERE operation_id = ?',
('2026-09-20T00:00:00+00:00', first_id),
)
self.db.conn.execute(
'UPDATE runtime_operations SET updated_at = ? WHERE operation_id = ?',
('2026-09-20T00:00:01+00:00', second_id),
)
self.db.conn.commit()
recent = self.db.recent_runtime_operations(1)
self.assertEqual(len(recent), 1)
self.assertEqual(recent[0]['operation_id'], second_id)
self.assertEqual(recent[0]['status'], 'requested')
for invalid in (True, 0, 501):
with self.subTest(limit=invalid), self.assertRaisesRegex(
ValueError, 'limit is out of range',
):
self.db.recent_runtime_operations(invalid)
def test_pending_host_agent_operations_are_bounded_and_prioritize_running(self):
requested_id, _ = self.create('apply-config')
running_id, _ = self.create('restart')
terminal_id, _ = self.create('restart')
self.db.mark_runtime_operation_running(running_id)
self.db.mark_runtime_operation_running(terminal_id)
self.db.reconcile_runtime_operation_result(self.envelope(
terminal_id, 'restart', 'succeeded',
identity={
'active_config_sha256': '1' * 64,
'active_secrets_sha256': '2' * 64,
},
))
pending = self.db.pending_runtime_agent_operations(limit=2)
self.assertEqual(
[item['operation_id'] for item in pending],
[running_id, requested_id],
)
self.assertEqual([item['status'] for item in pending], ['running', 'requested'])
self.assertNotIn(terminal_id, {item['operation_id'] for item in pending})
for invalid in (True, 0, 129):
with self.subTest(limit=invalid), self.assertRaisesRegex(
ValueError, 'limit is out of range',
):
self.db.pending_runtime_agent_operations(limit=invalid)
def test_audit_events_are_bounded_cursor_paginated_and_survive_reopen(self):
operation_ids = [self.create('apply-config')[0] for _ in range(4)]
self.db.mark_runtime_operation_running(operation_ids[-1])
resulting_identity = {
'active_config_sha256': '3' * 64,
'active_secrets_sha256': '2' * 64,
}
self.db.reconcile_runtime_operation_result(self.envelope(
operation_ids[-1], 'apply-config', 'succeeded',
identity=resulting_identity,
))
event_rows = self.db.conn.execute(
'SELECT id, operation_id FROM runtime_audit_events ORDER BY id DESC',
).fetchall()
first = self.db.runtime_audit_events(limit=3)
self.assertEqual(
[event['id'] for event in first['events']],
[event_rows[0]['id'], event_rows[1]['id'], event_rows[2]['id']],
)
self.assertEqual(
first['next_before_event_id'], event_rows[2]['id'],
)
self.assertEqual(
first['events'][0]['operation_id'], operation_ids[-1],
)
self.assertEqual(first['events'][0]['result'], 'succeeded')
self.assertEqual(
first['events'][0]['before_identity'], self.identity('apply-config'),
)
self.assertEqual(
first['events'][0]['after_identity'], resulting_identity,
)
self.db.close()
self.db = ScannerDB(db_path=self.path)
second = self.db.runtime_audit_events(
before_event_id=first['next_before_event_id'], limit=3,
)
self.assertEqual(
[event['id'] for event in second['events']],
[event_rows[3]['id'], event_rows[4]['id']],
)
self.assertIsNone(second['next_before_event_id'])
self.assertTrue(
set(event['id'] for event in first['events']).isdisjoint(
event['id'] for event in second['events']
)
)
def test_audit_event_page_rejects_unbounded_or_invalid_inputs(self):
for before_event_id in (True, 0, -1, 9223372036854775808):
with self.subTest(before_event_id=before_event_id), self.assertRaisesRegex(
ValueError, 'cursor is out of range',
):
self.db.runtime_audit_events(before_event_id=before_event_id)
for limit in (True, 0, 201):
with self.subTest(limit=limit), self.assertRaisesRegex(
ValueError, 'limit is out of range',
):
self.db.runtime_audit_events(limit=limit)
def test_recent_operations_deduplicate_status_transition_rows(self):
operation_id, _ = self.create('restart')
requested = self.db.conn.execute(
'SELECT * FROM runtime_operations WHERE operation_id = ?',
(operation_id,),
).fetchone()
self.db.mark_runtime_operation_running(operation_id)
running = self.db.conn.execute(
'SELECT * FROM runtime_operations WHERE operation_id = ?',
(operation_id,),
).fetchone()
original = self.db.conn
connection = mock.MagicMock()
connection.is_sqlite = False
def execute(sql, parameters=None):
if sql.startswith('SET TRANSACTION ISOLATION LEVEL'):
return mock.Mock()
status = parameters[0]
result = mock.Mock()
result.fetchall.return_value = {
'requested': [requested], 'running': [running],
}.get(status, [])
return result
connection.execute.side_effect = execute
self.db.conn = connection
try:
recent = self.db.recent_runtime_operations(10)
finally:
self.db.conn = original
self.assertEqual(len(recent), 1)
self.assertEqual(recent[0]['operation_id'], operation_id)
self.assertEqual(recent[0]['status'], 'running')
connection.commit.assert_called_once_with()
def test_requested_operations_reconcile_to_each_non_success_terminal_state(self):
cases = (
('failed', None, 'apply_failed'),
('failed_hold', None, 'rollback_failed'),
(
'rolled_back',
{
'active_config_sha256': '1' * 64,
'active_secrets_sha256': '2' * 64,
},
'health_check_failed',
),
)
for result, resulting, category in cases:
with self.subTest(result=result):
operation_id, _ = self.create('apply-both')
completed = self.db.reconcile_runtime_operation_result(self.envelope(
operation_id, 'apply-both', result,
identity=resulting, category=category,
detail='bounded_result',
))
self.assertEqual(completed['status'], result)
self.assertEqual(completed['agent_state'], result)
self.assertEqual(completed['safe_category'], category)
self.assertEqual(completed['safe_detail'], 'bounded_result')
self.assertEqual(completed['resulting_identity'], resulting)
with self.assertRaises(scanner_db.RuntimeOperationTransitionError):
self.db.mark_runtime_operation_running(operation_id)
def test_result_envelope_is_bounded_and_conflicting_replay_is_rejected(self):
operation_id, _ = self.create('restart')
resulting = {
'active_config_sha256': '1' * 64,
'active_secrets_sha256': '2' * 64,
}
envelope = self.envelope(
operation_id, 'restart', 'succeeded', identity=resulting,
)
self.db.reconcile_runtime_operation_result(envelope)
conflicting = self.envelope(
operation_id, 'restart', 'failed',
category='restart_failed', detail='bounded_result',
)
with self.assertRaises(scanner_db.RuntimeOperationIdentityConflictError):
self.db.reconcile_runtime_operation_result(conflicting)
with self.assertRaisesRegex(ValueError, 'byte bound'):
self.db.reconcile_runtime_operation_result(
envelope + (b' ' * scanner_db.RUNTIME_AGENT_RESULT_MAX_BYTES),
)
def test_success_identity_and_safe_codes_are_closed(self):
operation_id, _ = self.create('apply-secrets')
wrong_identity = {
'active_config_sha256': '1' * 64,
'active_secrets_sha256': '9' * 64,
}
with self.assertRaises(scanner_db.RuntimeOperationIdentityConflictError):
self.db.reconcile_runtime_operation_result(self.envelope(
operation_id, 'apply-secrets', 'succeeded', identity=wrong_identity,
))
with self.assertRaisesRegex(ValueError, 'safe category'):
self.db.reconcile_runtime_operation_result(self.envelope(
operation_id, 'apply-secrets', 'failed',
category='sk_live_aaaaaaaaaaaaaaaa', detail='bounded_result',
))
self.assertEqual(self.db.runtime_operation(operation_id)['status'], 'requested')
self.assertEqual(int(self.db.conn.execute(
'SELECT COUNT(*) AS count FROM runtime_audit_events WHERE operation_id = ?',
(operation_id,),
).fetchone()['count']), 1)
def test_deep_result_json_is_rejected_as_invalid(self):
payload = (b'[' * 4000) + (b']' * 4000)
with self.assertRaisesRegex(ValueError, 'invalid JSON'):
self.db.reconcile_runtime_operation_result(payload)
bracket_string = json.dumps(
{'brackets': '[\\"{' * 40},
sort_keys=True, separators=(',', ':'), ensure_ascii=True,
).encode('ascii')
with self.assertRaisesRegex(ValueError, 'shape is invalid'):
self.db.reconcile_runtime_operation_result(bracket_string)
def test_malformed_actions_are_rejected_without_durable_side_effects(self):
for action in (None, True, 1, [], {}, 'apply-config; restart'):
with self.subTest(action=action), self.assertRaisesRegex(
ValueError, 'action is invalid',
):
self.db.create_runtime_operation(
operation_id=str(uuid.uuid4()), actor='operator:alice',
action=action, expected_identity=self.identity('apply-config'),
)
counts = self.db.conn.execute(
'''SELECT
(SELECT COUNT(*) FROM runtime_operations) AS operations,
(SELECT COUNT(*) FROM runtime_audit_events) AS audits'''
).fetchone()
self.assertEqual((counts['operations'], counts['audits']), (0, 0))
def test_restart_boundaries_preserve_requested_running_and_reconciliation(self):
operation_id, created = self.create('apply-both')
self.assertEqual(created['status'], 'requested')
self.db.close()
self.db = ScannerDB(db_path=self.path)
running = self.db.mark_runtime_operation_running(operation_id)
self.assertEqual(running['status'], 'running')
self.db.close()
self.db = ScannerDB(db_path=self.path)
resulting = {
'active_config_sha256': '3' * 64,
'active_secrets_sha256': '4' * 64,
}
completed = self.db.reconcile_runtime_operation_result(self.envelope(
operation_id, 'apply-both', 'succeeded', identity=resulting,
))
self.assertEqual(completed['status'], 'succeeded')
self.assertEqual(completed['resulting_identity'], resulting)
self.db.close()
self.db = ScannerDB(db_path=self.path)
persisted = self.db.runtime_operation(operation_id)
self.assertEqual(persisted['status'], 'succeeded')
self.assertEqual(persisted['resulting_identity'], resulting)
def test_audit_records_are_content_free_for_hostile_result_fields(self):
operation_id, _ = self.create('restart')
sentinels = (
'authorization-bearer-must-not-persist',
'csrf-must-not-persist',
'candidate-content-must-not-persist',
'stdout-must-not-persist',
)
hostile = {
'schema': 1,
'operation_id': operation_id,
'action': 'restart',
'result': 'succeeded',
'safe_category': None,
'safe_detail': None,
'resulting_identity': {
'active_config_sha256': '1' * 64,
'active_secrets_sha256': '2' * 64,
},
'authorization': sentinels[0],
'csrf': sentinels[1],
'content': sentinels[2],
'stdout': sentinels[3],
}
with self.assertRaisesRegex(ValueError, 'shape'):
self.db.reconcile_runtime_operation_result(json.dumps(hostile).encode('utf-8'))
completed = self.db.reconcile_runtime_operation_result(self.envelope(
operation_id,
'restart',
'succeeded',
identity={
'active_config_sha256': '1' * 64,
'active_secrets_sha256': '2' * 64,
},
))
self.assertEqual(completed['status'], 'succeeded')
operation = self.db.conn.execute(
'SELECT * FROM runtime_operations WHERE operation_id = ?',
(operation_id,),
).fetchone()
audits = self.db.conn.execute(
'SELECT * FROM runtime_audit_events WHERE operation_id = ? ORDER BY id',
(operation_id,),
).fetchall()
serialized = json.dumps({
'operation': dict(operation),
'audits': [dict(row) for row in audits],
}, sort_keys=True, default=str)
for sentinel in sentinels:
self.assertNotIn(sentinel, serialized)
self.assertEqual([row['result'] for row in audits], ['accepted', 'succeeded'])
def test_managed_source_operations_are_attributed_audited_and_replay_safe(self):
operation_id = str(uuid.uuid4())
created = self.db.create_runtime_source_operation(
operation_id=operation_id, actor='Alice.Operator',
source_id='discovery-producer:gitlab', source_action='set-interval',
interval_seconds=3600,
)
replay = self.db.create_runtime_source_operation(
operation_id=operation_id, actor='Alice.Operator',
source_id='discovery-producer:gitlab', source_action='set-interval',
interval_seconds=3600,
)
completed = self.db.complete_runtime_source_operation(
operation_id, succeeded=True, outcome='completed',
)
completed_replay = self.db.complete_runtime_source_operation(
operation_id, succeeded=True, outcome='completed',
)
self.assertEqual(created['status'], 'running')
self.assertFalse(created['replayed'])
self.assertTrue(replay['replayed'])
self.assertEqual(completed['status'], 'succeeded')
self.assertEqual(completed['agent_state'], 'not_required')
self.assertEqual(completed['expected_identity'], {
'source_id': 'discovery-producer:gitlab',
'source_action': 'set-interval',
'interval_seconds': 3600,
})
self.assertEqual(completed['resulting_identity'], {
'source_id': 'discovery-producer:gitlab',
'source_action': 'set-interval',
'outcome': 'completed',
})
self.assertTrue(completed_replay['replayed'])
events = self.db.conn.execute(
'SELECT * FROM runtime_audit_events WHERE operation_id = ? ORDER BY id',
(operation_id,),
).fetchall()
self.assertEqual([row['result'] for row in events], ['accepted', 'succeeded'])
self.assertTrue(all(row['actor'] == 'Alice.Operator' for row in events))
self.assertTrue(all(row['target_kind'] == 'managed-source' for row in events))
self.assertTrue(all(
row['target_ref'] == 'discovery-producer:gitlab' for row in events
))
failed_id = str(uuid.uuid4())
self.db.create_runtime_source_operation(
operation_id=failed_id, actor='Alice.Operator',
source_id='discovery-producer:huggingface', source_action='restart',
)
failed = self.db.complete_runtime_source_operation(
failed_id, succeeded=False,
)
self.assertEqual(failed['status'], 'failed')
self.assertEqual(failed['safe_category'], 'supervisor_action_failed')
failed_events = self.db.conn.execute(
'SELECT result, safe_category FROM runtime_audit_events '
'WHERE operation_id = ? ORDER BY id', (failed_id,),
).fetchall()
self.assertEqual(
[(row['result'], row['safe_category']) for row in failed_events],
[('accepted', None), ('failed', 'supervisor_action_failed')],
)
extended_cases = (
('dashboard', 'restart', {}, {
'source_id': 'dashboard', 'source_action': 'restart',
'interval_seconds': None,
}),
('keychecks', 'set-mode', {'mode': 'repeat'}, {
'source_id': 'keychecks', 'source_action': 'set-mode',
'interval_seconds': None, 'mode': 'repeat',
'restart_enabled': None, 'restart_delay_seconds': None,
}),
('result-ingester', 'set-restart', {'restart_enabled': False}, {
'source_id': 'result-ingester', 'source_action': 'set-restart',
'interval_seconds': None, 'mode': None,
'restart_enabled': False, 'restart_delay_seconds': None,
}),
('jsonl-projector', 'set-restart-delay', {
'restart_delay_seconds': 30,
}, {
'source_id': 'jsonl-projector',
'source_action': 'set-restart-delay', 'interval_seconds': None,
'mode': None, 'restart_enabled': None,
'restart_delay_seconds': 30,
}),
)
for source_id, source_action, parameters, expected in extended_cases:
with self.subTest(source_id=source_id, source_action=source_action):
state = self.db.create_runtime_source_operation(
operation_id=str(uuid.uuid4()), actor='Alice.Operator',
source_id=source_id, source_action=source_action, **parameters,
)
self.assertEqual(state['expected_identity'], expected)
for invalid in (
{'source_id': '../discovery-producer:github', 'source_action': 'start'},
{'source_id': 'discovery-producer:gitlab', 'source_action': 'run'},
):
with self.subTest(invalid=invalid), self.assertRaisesRegex(
ValueError, 'request is invalid',
):
self.db.create_runtime_source_operation(
operation_id=str(uuid.uuid4()), actor='Alice.Operator',
**invalid,
)
def test_worker_admin_operations_are_hash_bound_attributed_and_content_free(self):
operation_id = str(uuid.uuid4())
created = self.db.create_runtime_worker_admin_operation(
operation_id=operation_id, actor='Alice.Operator',
action='workers.device.issue', target_ref='device-a',
request_sha256='7' * 64,
)
replay = self.db.create_runtime_worker_admin_operation(
operation_id=operation_id, actor='Alice.Operator',
action='workers.device.issue', target_ref='device-a',
request_sha256='7' * 64,
)
completed = self.db.complete_runtime_worker_admin_operation(
operation_id, succeeded=True, affected_count=1,
)
completed_replay = self.db.complete_runtime_worker_admin_operation(
operation_id, succeeded=True, affected_count=1,
)
self.assertFalse(created['replayed'])
self.assertTrue(replay['replayed'])
self.assertEqual(completed['status'], 'succeeded')
self.assertTrue(completed_replay['replayed'])
self.assertEqual(completed['expected_identity'], {
'action': 'workers.device.issue',
'target_ref': 'device-a',
'request_sha256': '7' * 64,
})
self.assertEqual(completed['resulting_identity'], {
'action': 'workers.device.issue',
'target_ref': 'device-a',
'outcome': 'completed',
'affected_count': 1,
})
events = self.db.conn.execute(
'SELECT * FROM runtime_audit_events WHERE operation_id = ? ORDER BY id',
(operation_id,),
).fetchall()
self.assertEqual([row['result'] for row in events], ['accepted', 'succeeded'])
self.assertTrue(all(row['actor'] == 'Alice.Operator' for row in events))
self.assertTrue(all(row['target_kind'] == 'worker-admin' for row in events))
serialized = json.dumps([dict(row) for row in events], sort_keys=True)
self.assertNotIn('one-time-device-token', serialized)
with self.assertRaises(scanner_db.RuntimeOperationIdentityConflictError):
self.db.create_runtime_worker_admin_operation(
operation_id=operation_id, actor='Alice.Operator',
action='workers.device.issue', target_ref='device-a',
request_sha256='8' * 64,
)
failed_id = str(uuid.uuid4())
self.db.create_runtime_worker_admin_operation(
operation_id=failed_id, actor='Alice.Operator',
action='workers.user.disable', target_ref='alice',
request_sha256='9' * 64,
)
failed = self.db.complete_runtime_worker_admin_operation(
failed_id, succeeded=False,
)
self.assertEqual(failed['status'], 'failed')
self.assertEqual(failed['safe_category'], 'worker_admin_mutation_failed')
def test_runtime_document_save_operation_is_hash_only_and_replay_safe(self):
operation_id = str(uuid.uuid4())
arguments = {
'operation_id': operation_id,
'actor': 'Alice.Operator',
'action': 'runtime.secrets.save',
'active_config_sha256': '1' * 64,
'active_secrets_sha256': '2' * 64,
'candidate_config_sha256': '3' * 64,
'candidate_secrets_sha256': '4' * 64,
'candidate_after_sha256': '5' * 64,
'candidate_before_bytes': 123,
'candidate_after_bytes': 456,
'candidate_before_present': True,
}
created = self.db.create_runtime_document_operation(**arguments)
replay = self.db.create_runtime_document_operation(**arguments)
completed = self.db.complete_runtime_document_operation(
operation_id, succeeded=True, candidate_sha256='5' * 64,
written=True,
)
terminal_replay = self.db.complete_runtime_document_operation(
operation_id, succeeded=True, candidate_sha256='5' * 64,
written=True,
)
self.assertFalse(created['replayed'])
self.assertTrue(replay['replayed'])
self.assertEqual(completed['status'], 'succeeded')
self.assertTrue(terminal_replay['replayed'])
self.assertEqual(completed['target_kind'], 'runtime-document')
self.assertEqual(completed['expected_identity'], {
'document': 'secrets',
'active_config_sha256': '1' * 64,
'active_secrets_sha256': '2' * 64,
'candidate_config_sha256': '3' * 64,
'candidate_secrets_sha256': '4' * 64,
'candidate_after_sha256': '5' * 64,
'candidate_before_bytes': 123,
'candidate_after_bytes': 456,
'candidate_before_present': True,
})
self.assertEqual(completed['resulting_identity'], {
'document': 'secrets', 'outcome': 'completed',
'candidate_sha256': '5' * 64, 'written': True,
})
rows = self.db.conn.execute(
'SELECT * FROM runtime_audit_events WHERE operation_id = ? ORDER BY id',
(operation_id,),
).fetchall()
serialized = json.dumps([dict(row) for row in rows], sort_keys=True)
self.assertEqual([row['result'] for row in rows], ['accepted', 'succeeded'])
self.assertNotIn('plaintext-secret', serialized)
self.assertEqual(int(rows[0]['before_bytes']), 123)
self.assertEqual(int(rows[0]['after_bytes']), 456)
self.assertEqual(int(rows[1]['before_bytes']), 123)
self.assertEqual(int(rows[1]['after_bytes']), 456)
with self.assertRaises(scanner_db.RuntimeOperationIdentityConflictError):
self.db.complete_runtime_document_operation(
operation_id, succeeded=True, candidate_sha256='6' * 64,
written=True,
)
conflicts = {
'candidate_config_sha256': '6' * 64,
'candidate_secrets_sha256': '7' * 64,
'candidate_after_sha256': '8' * 64,
'candidate_before_bytes': 124,
'candidate_after_bytes': 457,
'candidate_before_present': False,
}
for name, value in conflicts.items():
with self.subTest(name=name):
with self.assertRaises(scanner_db.RuntimeOperationIdentityConflictError):
self.db.create_runtime_document_operation(
**{**arguments, name: value},
)
def test_managed_file_operations_are_content_free_and_replay_safe(self):
secret = 'uploaded-file-content-must-not-persist'
cases = (
(
'files.create', None, '2' * 64, 12,
{
'before_sha256': None, 'before_byte_count': None,
'after_sha256': '2' * 64, 'after_byte_count': 12,
'written': True,
},
),
(
'files.replace', '1' * 64, '2' * 64, 12,
{
'before_sha256': '1' * 64, 'before_byte_count': 10,
'after_sha256': '2' * 64, 'after_byte_count': 12,
'written': True,
},
),
(
'files.delete', '1' * 64, None, None,
{
'before_sha256': '1' * 64, 'before_byte_count': 10,
'after_sha256': None, 'after_byte_count': None,
'written': True,
},
),
)
for action, expected_hash, proposed_hash, proposed_bytes, result in cases:
with self.subTest(action=action):
operation_id = str(uuid.uuid4())
arguments = {
'operation_id': operation_id,
'actor': 'Alice.Operator',
'action': action,
'root_id': 'exports',
'relative_path': 'nested/result.bin',
'expected_sha256': expected_hash,
'proposed_sha256': proposed_hash,
'proposed_byte_count': proposed_bytes,
}
created = self.db.create_runtime_managed_file_operation(**arguments)
replay = self.db.create_runtime_managed_file_operation(**arguments)
completed = self.db.complete_runtime_managed_file_operation(
operation_id, succeeded=True, **result,
)
terminal_replay = self.db.complete_runtime_managed_file_operation(
operation_id, succeeded=True, **result,
)
self.assertFalse(created['replayed'])
self.assertTrue(replay['replayed'])
self.assertEqual(completed['status'], 'succeeded')
self.assertTrue(terminal_replay['replayed'])
self.assertEqual(completed['target_kind'], 'managed-file')
self.assertEqual(completed['target_ref'], 'exports')
self.assertEqual(completed['expected_identity'], {
'root_id': 'exports',
'relative_path': 'nested/result.bin',
'expected_sha256': expected_hash,
'proposed_sha256': proposed_hash,
'proposed_byte_count': proposed_bytes,
})
self.assertEqual(completed['resulting_identity'], {
'root_id': 'exports',
'relative_path': 'nested/result.bin',
'outcome': 'completed',
**result,
})
rows = self.db.conn.execute(
'SELECT * FROM runtime_audit_events WHERE operation_id = ? ORDER BY id',
(operation_id,),
).fetchall()
self.assertEqual(
[row['result'] for row in rows], ['accepted', 'succeeded'],
)
self.assertNotIn(
secret, json.dumps([dict(row) for row in rows], sort_keys=True),
)
for name, value in (
('actor', 'Mallory'),
('relative_path', 'other.bin'),
('expected_sha256', None if expected_hash else '1' * 64),
('proposed_sha256', '3' * 64 if proposed_hash else '3' * 64),
('proposed_byte_count', 13 if proposed_bytes is not None else 13),
):
conflicting = {**arguments, name: value}
with self.assertRaises((
ValueError, scanner_db.RuntimeOperationIdentityConflictError,
)):
self.db.create_runtime_managed_file_operation(**conflicting)
failed_id = str(uuid.uuid4())
self.db.create_runtime_managed_file_operation(
operation_id=failed_id, actor='Alice.Operator', action='files.delete',
root_id='exports', relative_path='failed.bin',
expected_sha256='4' * 64, proposed_sha256=None,
proposed_byte_count=None,
)
failed = self.db.complete_runtime_managed_file_operation(
failed_id, succeeded=False,
)
self.assertEqual(failed['status'], 'failed')
self.assertEqual(failed['safe_category'], 'managed_file_mutation_failed')
def test_managed_file_execution_lock_serializes_sqlite_connections(self):
operation_id = str(uuid.uuid4())
started = threading.Event()
acquired = threading.Event()
release_contender = threading.Event()
errors = []
self.db.acquire_runtime_managed_file_execution(operation_id)
def contend():
contender = ScannerDB(db_path=self.path, initialize=False)
try:
started.set()
contender.acquire_runtime_managed_file_execution(operation_id)
acquired.set()
if not release_contender.wait(2):
raise AssertionError('managed file execution lock release timed out')
contender.release_runtime_managed_file_execution(operation_id)
except BaseException as exc:
errors.append(exc)
finally:
contender.close()
thread = threading.Thread(target=contend, daemon=True)
thread.start()
first_released = False
try:
self.assertTrue(started.wait(2))
self.assertFalse(acquired.wait(0.1))
self.db.release_runtime_managed_file_execution(operation_id)
first_released = True
self.assertTrue(acquired.wait(2))
finally:
if not first_released:
try:
self.db.release_runtime_managed_file_execution(operation_id)
except scanner_db.RuntimeOperationTransitionError:
pass
release_contender.set()
thread.join(2)
self.assertFalse(thread.is_alive())
self.assertEqual(errors, [])
def test_managed_file_execution_close_finishes_cleanup_on_cancellation(self):
class CancelLock:
def release(self):
raise asyncio.CancelledError()
class Connection:
def __init__(self):
self.closed = False
def close(self):
self.closed = True
db = ScannerDB(enabled=False, db_url='')
connection = Connection()
db.conn = connection
db._runtime_managed_file_execution_held = (
str(uuid.uuid4()), 1, CancelLock(),
)
with self.assertRaises(asyncio.CancelledError):
db.close()
self.assertTrue(connection.closed)
self.assertIsNone(db.conn)
self.assertIsNone(db._runtime_managed_file_execution_held)
if __name__ == '__main__':
unittest.main()