Add continuous material polling runner
This commit is contained in:
@@ -0,0 +1,277 @@
|
||||
import io
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from pathlib import Path
|
||||
from unittest.mock import Mock, patch
|
||||
|
||||
import pytest
|
||||
import yaml
|
||||
|
||||
from production_analytics.calculations.config import (
|
||||
CalculationConfigError,
|
||||
load_material_calculation,
|
||||
)
|
||||
from production_analytics.calculations.material_consumption import MaterialIntegrationState
|
||||
from production_analytics.cli.__main__ import main
|
||||
from production_analytics.enlyze.exploration import ConfigurationError
|
||||
from production_analytics.enlyze.gateway import EnlyzeProductionRun
|
||||
from production_analytics.service.material_polling import MaterialPollResult
|
||||
from production_analytics.service.material_runner import MaterialPollingRunner, MaterialStateError
|
||||
from production_analytics.service.material_runtime import (
|
||||
_ReportingStateStore,
|
||||
build_material_runner,
|
||||
)
|
||||
|
||||
EXAMPLE = Path(__file__).resolve().parents[1] / 'config/k7-material-consumption.yaml'
|
||||
NOW = datetime(2026, 1, 1, tzinfo=UTC)
|
||||
|
||||
|
||||
def write_config(tmp_path, entries):
|
||||
path = tmp_path / 'calculations.yaml'
|
||||
path.write_text(yaml.safe_dump({'calculations': entries}))
|
||||
return path
|
||||
|
||||
|
||||
def entry():
|
||||
return yaml.safe_load(EXAMPLE.read_text())['calculations'][0]
|
||||
|
||||
|
||||
def test_valid_k7_config():
|
||||
config = load_material_calculation(EXAMPLE)
|
||||
assert config.id == 'k7-fiber-consumption'
|
||||
assert config.machine_ref == 'c220f95c-a65e-4cb7-99b7-0626d6c7508c'
|
||||
assert config.rate_signal_ref == 'c9d06af5-f6d6-4ede-b6c4-5a98bac77129'
|
||||
assert config.gate_signal_ref == '823867bb-f5d2-40eb-b875-657155addfd0'
|
||||
assert config.gate_threshold == 0.5
|
||||
assert config.max_sample_gap_seconds == 20.0
|
||||
|
||||
|
||||
@pytest.mark.parametrize('field', list(entry()))
|
||||
def test_missing_fields(tmp_path, field):
|
||||
data = entry()
|
||||
del data[field]
|
||||
with pytest.raises(CalculationConfigError, match=f'missing fields: {field}'):
|
||||
load_material_calculation(write_config(tmp_path, [data]))
|
||||
|
||||
|
||||
@pytest.mark.parametrize(('field', 'value'), [
|
||||
('type', 'integration'), ('version', 1), ('version', '2'),
|
||||
('machine_ref', ''), ('rate_signal_ref', ' '), ('gate_signal_ref', None), ('id', ''),
|
||||
('gate_threshold', True), ('gate_threshold', '0.5'), ('gate_threshold', float('nan')),
|
||||
('max_sample_gap_seconds', float('inf')), ('max_sample_gap_seconds', 0),
|
||||
('max_sample_gap_seconds', -1), ('output_unit', 'tonnes'), ('group_by', 'run'),
|
||||
('gate_treshold', 1),
|
||||
])
|
||||
def test_invalid_fields(tmp_path, field, value):
|
||||
data = entry()
|
||||
data[field] = value
|
||||
with pytest.raises(CalculationConfigError):
|
||||
load_material_calculation(write_config(tmp_path, [data]))
|
||||
|
||||
|
||||
@pytest.mark.parametrize('content', [
|
||||
'calculations: [', 'calculations: []', 'calculations: {}', 'other: []',
|
||||
'calculations: []\ncalculations: []', 'calculations:\n - id: x\n id: y',
|
||||
])
|
||||
def test_malformed_documents(tmp_path, content):
|
||||
path = tmp_path / 'bad.yaml'
|
||||
path.write_text(content)
|
||||
with pytest.raises(CalculationConfigError):
|
||||
load_material_calculation(path)
|
||||
|
||||
|
||||
def test_selection(tmp_path):
|
||||
first, second = entry(), dict(entry(), id='second')
|
||||
path = write_config(tmp_path, [first, second])
|
||||
with pytest.raises(CalculationConfigError, match='specify --calculation-id'):
|
||||
load_material_calculation(path)
|
||||
assert load_material_calculation(path, 'second').id == 'second'
|
||||
with pytest.raises(CalculationConfigError, match='not found'):
|
||||
load_material_calculation(path, 'absent')
|
||||
with pytest.raises(CalculationConfigError, match='unique'):
|
||||
load_material_calculation(write_config(tmp_path, [first, first]), first['id'])
|
||||
|
||||
|
||||
def test_repeated_sequential_polls_and_long_cycle():
|
||||
events = []
|
||||
elapsed = 0
|
||||
|
||||
def poll(*, now):
|
||||
nonlocal elapsed
|
||||
assert now == NOW + timedelta(seconds=elapsed)
|
||||
events.append('poll-start')
|
||||
elapsed += 25 # Longer than the interval; no real sleep or concurrent work.
|
||||
events.append('poll-end')
|
||||
|
||||
def sleep(seconds):
|
||||
nonlocal elapsed
|
||||
events.append(('sleep', seconds))
|
||||
elapsed += seconds
|
||||
if elapsed >= 70:
|
||||
raise KeyboardInterrupt
|
||||
|
||||
output = io.StringIO()
|
||||
service = Mock(poll_once=Mock(side_effect=poll))
|
||||
MaterialPollingRunner(
|
||||
service, machine_id='machine', poll_interval_seconds=10,
|
||||
clock=lambda: NOW + timedelta(seconds=elapsed), sleep=sleep, stdout=output,
|
||||
).run()
|
||||
assert events == ['poll-start', 'poll-end', ('sleep', 10)] * 2
|
||||
assert service.poll_once.call_count == 2
|
||||
assert output.getvalue().count('no open Production Run') == 2
|
||||
assert 'stopped' in output.getvalue()
|
||||
|
||||
|
||||
def test_error_then_success_preserves_opaque_order_and_reports_totals():
|
||||
run = EnlyzeProductionRun('run', 'machine', None, ' 00842/ABC ', NOW, None)
|
||||
result = MaterialPollResult(run, MaterialIntegrationState(12.5, 20))
|
||||
service = Mock(poll_once=Mock(side_effect=[ValueError('SECRET'), result]))
|
||||
output, errors = io.StringIO(), io.StringIO()
|
||||
sleep = Mock(side_effect=[None, KeyboardInterrupt])
|
||||
MaterialPollingRunner(
|
||||
service, machine_id='machine', poll_interval_seconds=5, clock=lambda: NOW,
|
||||
sleep=sleep, stdout=output, stderr=errors,
|
||||
).run()
|
||||
assert service.poll_once.call_count == 2
|
||||
assert sleep.call_count == 2
|
||||
assert 'gateway/poll cycle: ValueError' in errors.getvalue()
|
||||
assert 'SECRET' not in errors.getvalue()
|
||||
assert "production_order=' 00842/ABC '" in output.getvalue()
|
||||
assert 'consumption_kg=12.500000000 running_seconds=20.000' in output.getvalue()
|
||||
assert NOW.isoformat() in output.getvalue()
|
||||
|
||||
|
||||
def test_interrupt_during_poll():
|
||||
sleep, output = Mock(), io.StringIO()
|
||||
MaterialPollingRunner(
|
||||
Mock(poll_once=Mock(side_effect=KeyboardInterrupt)), machine_id='m',
|
||||
poll_interval_seconds=1, sleep=sleep, stdout=output,
|
||||
).run()
|
||||
sleep.assert_not_called()
|
||||
assert 'stopped' in output.getvalue()
|
||||
|
||||
|
||||
def test_cycle_configuration_failure_is_fatal():
|
||||
sleep = Mock()
|
||||
runner = MaterialPollingRunner(
|
||||
Mock(poll_once=Mock(side_effect=ConfigurationError('bad settings'))),
|
||||
machine_id='m', poll_interval_seconds=1, sleep=sleep,
|
||||
)
|
||||
with pytest.raises(ConfigurationError):
|
||||
runner.run()
|
||||
sleep.assert_not_called()
|
||||
|
||||
|
||||
@pytest.mark.parametrize('interval', [0, -1, float('nan'), float('inf'), True])
|
||||
def test_invalid_interval(interval):
|
||||
with pytest.raises(CalculationConfigError):
|
||||
MaterialPollingRunner(Mock(), machine_id='m', poll_interval_seconds=interval)
|
||||
|
||||
|
||||
@pytest.mark.parametrize('operation', ['load', 'save'])
|
||||
def test_state_error_classification(operation):
|
||||
store = Mock()
|
||||
getattr(store, operation).side_effect = OSError('SECRET')
|
||||
wrapped = _ReportingStateStore(store)
|
||||
arguments = ('machine', ' 00842 ') if operation == 'load' else ('machine', ' 00842 ', Mock())
|
||||
phase = 'loading' if operation == 'load' else 'saving'
|
||||
with pytest.raises(MaterialStateError, match=f'state {phase}: OSError') as error:
|
||||
getattr(wrapped, operation)(*arguments)
|
||||
assert 'SECRET' not in str(error.value)
|
||||
|
||||
|
||||
def test_runtime_wiring(tmp_path, monkeypatch):
|
||||
monkeypatch.setenv('ENLYZE_BASE_URL', 'https://example.invalid/api/')
|
||||
secrets = tmp_path / 'secret.env'
|
||||
secrets.write_text('ENLYZE_API_KEY=secret-from-file\n')
|
||||
with patch('production_analytics.service.material_runtime.MaterialPollingService') as service:
|
||||
runner = build_material_runner(
|
||||
load_material_calculation(EXAMPLE), poll_interval_seconds=7,
|
||||
state_directory=tmp_path / 'state', secrets_file=secrets,
|
||||
)
|
||||
kwargs = service.call_args.kwargs
|
||||
assert kwargs['machine_id'] == entry()['machine_ref']
|
||||
assert kwargs['rate_variable_id'] == entry()['rate_signal_ref']
|
||||
assert kwargs['gate_variable_id'] == entry()['gate_signal_ref']
|
||||
assert kwargs['gate_threshold'] == 0.5
|
||||
assert kwargs['max_sample_gap_seconds'] == 20.0
|
||||
assert kwargs['gateway']._client._settings.api_key == 'secret-from-file'
|
||||
assert kwargs['state_store'].store._directory == tmp_path / 'state'
|
||||
assert runner.service is service.return_value
|
||||
assert runner.interval == 7
|
||||
|
||||
|
||||
def test_cli_bad_config(tmp_path, capsys):
|
||||
assert main(['run', 'material-poll', '--config', str(tmp_path / 'missing')]) == 2
|
||||
assert 'Config/startup' in capsys.readouterr().err
|
||||
|
||||
|
||||
def test_cli_startup_and_success(tmp_path, capsys):
|
||||
args = ['run', 'material-poll', '--config', str(EXAMPLE),
|
||||
'--calculation-id', 'k7-fiber-consumption']
|
||||
with patch('production_analytics.service.material_runtime.build_material_runner') as build:
|
||||
assert main(args) == 0
|
||||
build.return_value.run.assert_called_once_with()
|
||||
assert build.call_args.kwargs['state_directory'] == Path('data/state/material')
|
||||
build.side_effect = ConfigurationError('bad settings')
|
||||
assert main(args) == 2
|
||||
assert 'bad settings' in capsys.readouterr().err
|
||||
|
||||
|
||||
def test_failed_cycle_keeps_checkpoint_and_recovers(tmp_path):
|
||||
from production_analytics.calculations.material_consumption import MaterialSample
|
||||
from production_analytics.service.material_polling import (
|
||||
MaterialPollingService,
|
||||
MaterialPollingState,
|
||||
)
|
||||
from production_analytics.service.material_state_store import JsonMaterialStateStore
|
||||
|
||||
order = ' 00842 '
|
||||
store = JsonMaterialStateStore(tmp_path)
|
||||
store.save('machine', order, MaterialPollingState('run', MaterialIntegrationState(
|
||||
cumulative_consumption_kg=1, integrated_running_seconds=10,
|
||||
last_processed_timestamp=NOW, last_material_rate_kg_per_hour=3600,
|
||||
last_gate_value=1, integration_active=True,
|
||||
)))
|
||||
checkpoint = next(tmp_path.glob('*.json'))
|
||||
original = checkpoint.read_bytes()
|
||||
gateway = Mock()
|
||||
gateway.get_open_production_run.return_value = EnlyzeProductionRun(
|
||||
'run', 'machine', None, order, NOW, None,
|
||||
)
|
||||
gateway.get_material_samples.side_effect = [
|
||||
ValueError('bad response'), [MaterialSample(NOW + timedelta(seconds=10), 3600, 1)],
|
||||
]
|
||||
service = MaterialPollingService(
|
||||
gateway=gateway, state_store=_ReportingStateStore(store), machine_id='machine',
|
||||
rate_variable_id='rate', gate_variable_id='gate', gate_threshold=.5,
|
||||
max_sample_gap_seconds=20,
|
||||
)
|
||||
sleeps = 0
|
||||
|
||||
def sleep(seconds):
|
||||
nonlocal sleeps
|
||||
sleeps += 1
|
||||
if sleeps == 1:
|
||||
assert checkpoint.read_bytes() == original
|
||||
else:
|
||||
raise KeyboardInterrupt
|
||||
|
||||
MaterialPollingRunner(
|
||||
service, machine_id='machine', poll_interval_seconds=1,
|
||||
clock=lambda: NOW + timedelta(seconds=10), sleep=sleep,
|
||||
stdout=io.StringIO(), stderr=io.StringIO(),
|
||||
).run()
|
||||
restored = store.load('machine', order)
|
||||
assert restored.integration_state.cumulative_consumption_kg == 11
|
||||
assert restored.integration_state.integrated_running_seconds == 20
|
||||
|
||||
|
||||
def test_unwritable_state_location_is_startup_failure(tmp_path, monkeypatch, capsys):
|
||||
monkeypatch.setenv('ENLYZE_BASE_URL', 'https://example.invalid/api/')
|
||||
blocked = tmp_path / 'file'
|
||||
blocked.write_text('not a directory')
|
||||
assert main([
|
||||
'run', 'material-poll', '--config', str(EXAMPLE),
|
||||
'--state-directory', str(blocked), '--secrets-file', str(tmp_path / 'absent'),
|
||||
]) == 2
|
||||
assert 'State directory startup check failed' in capsys.readouterr().err
|
||||
Reference in New Issue
Block a user