1053 lines
44 KiB
Python
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()
|