Files
production-analytics/tests/test_material_polling_service.py

410 lines
14 KiB
Python

from datetime import UTC, datetime
from unittest.mock import Mock
import pytest
from production_analytics.calculations.material_consumption import (
MaterialIntegrationState,
MaterialSample,
)
from production_analytics.enlyze.gateway import EnlyzeProductionRun
from production_analytics.service.material_polling import (
MaterialPollingService,
MaterialPollingState,
)
def test_poll_once_starts_at_run_start_without_saved_state() -> None:
gateway = Mock()
gateway.get_production_runs.return_value = []
state_store = Mock()
state_store.load.return_value = None
run = EnlyzeProductionRun(
uuid="run-1",
machine_id="machine-1",
product_id="product-1",
production_order="ORDER-1",
start=datetime(2026, 9, 3, 5, 7, 47, tzinfo=UTC),
end=None,
)
gateway.get_open_production_run.return_value = run
gateway.get_production_runs.return_value = [run]
samples = [
MaterialSample(
timestamp=datetime(2026, 9, 3, 5, 7, 50, tzinfo=UTC),
material_rate_kg_per_hour=3600.0,
gate_value=1.0,
),
MaterialSample(
timestamp=datetime(2026, 9, 3, 5, 8, 0, tzinfo=UTC),
material_rate_kg_per_hour=3600.0,
gate_value=1.0,
),
]
gateway.get_material_samples.return_value = samples
service = MaterialPollingService(
gateway=gateway,
state_store=state_store,
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
gate_threshold=0.5,
max_sample_gap_seconds=20.0,
)
result = service.poll_once(
now=datetime(2026, 9, 3, 5, 8, 10, tzinfo=UTC),
)
assert result is not None
assert result.run == run
assert result.state.cumulative_consumption_kg == 10.0
assert result.state.integrated_running_seconds == 10.0
gateway.get_material_samples.assert_called_once_with(
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
start=run.start,
end=datetime(2026, 9, 3, 5, 8, 10, tzinfo=UTC),
)
state_store.save.assert_called_once_with(
"machine-1",
run.production_order,
MaterialPollingState(
run_id=run.uuid,
integration_state=result.state,
),
)
def test_poll_once_resumes_from_saved_timestamp_without_double_counting() -> None:
gateway = Mock()
gateway.get_production_runs.return_value = []
state_store = Mock()
run = EnlyzeProductionRun(
uuid="run-1",
machine_id="machine-1",
product_id="product-1",
production_order="ORDER-1",
start=datetime(2026, 9, 3, 5, 7, 47, tzinfo=UTC),
end=None,
)
gateway.get_open_production_run.return_value = run
saved_state = MaterialIntegrationState(
cumulative_consumption_kg=10.0,
integrated_running_seconds=10.0,
last_processed_timestamp=datetime(2026, 9, 3, 5, 8, 0, tzinfo=UTC),
last_material_rate_kg_per_hour=3600.0,
last_gate_value=1.0,
integration_active=True,
)
state_store.load.return_value = MaterialPollingState(
run_id=run.uuid,
integration_state=saved_state,
)
gateway.get_material_samples.return_value = [
MaterialSample(
timestamp=datetime(2026, 9, 3, 5, 8, 0, tzinfo=UTC),
material_rate_kg_per_hour=3600.0,
gate_value=1.0,
),
MaterialSample(
timestamp=datetime(2026, 9, 3, 5, 8, 10, tzinfo=UTC),
material_rate_kg_per_hour=3600.0,
gate_value=1.0,
),
]
service = MaterialPollingService(
gateway=gateway,
state_store=state_store,
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
gate_threshold=0.5,
max_sample_gap_seconds=20.0,
)
result = service.poll_once(
now=datetime(2026, 9, 3, 5, 8, 20, tzinfo=UTC),
)
assert result is not None
gateway.get_production_runs.assert_not_called()
assert result.state.cumulative_consumption_kg == 20.0
assert result.state.integrated_running_seconds == 20.0
gateway.get_material_samples.assert_called_once_with(
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
start=saved_state.last_processed_timestamp,
end=datetime(2026, 9, 3, 5, 8, 20, tzinfo=UTC),
)
def test_poll_once_without_open_run_does_nothing() -> None:
gateway = Mock()
gateway.get_production_runs.return_value = []
state_store = Mock()
gateway.get_open_production_run.return_value = None
service = MaterialPollingService(
gateway=gateway,
state_store=state_store,
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
gate_threshold=0.5,
max_sample_gap_seconds=20.0,
)
result = service.poll_once(
now=datetime(2026, 9, 3, 5, 8, 20, tzinfo=UTC),
)
assert result is None
gateway.get_material_samples.assert_not_called()
state_store.load.assert_not_called()
state_store.save.assert_not_called()
def test_new_run_of_same_order_keeps_total_but_restarts_integration_at_run_start() -> None:
gateway = Mock()
gateway.get_production_runs.return_value = []
state_store = Mock()
run = EnlyzeProductionRun(
uuid="run-2",
machine_id="machine-1",
product_id="product-1",
production_order="ORDER-1",
start=datetime(2026, 9, 3, 8, 0, 0, tzinfo=UTC),
end=None,
)
gateway.get_open_production_run.return_value = run
previous_integration_state = MaterialIntegrationState(
cumulative_consumption_kg=100.0,
integrated_running_seconds=360.0,
last_processed_timestamp=datetime(2026, 9, 3, 6, 0, 0, tzinfo=UTC),
last_material_rate_kg_per_hour=3600.0,
last_gate_value=1.0,
integration_active=True,
)
state_store.load.return_value = MaterialPollingState(
run_id="run-1",
integration_state=previous_integration_state,
)
gateway.get_material_samples.return_value = [
MaterialSample(
timestamp=datetime(2026, 9, 3, 8, 0, 10, tzinfo=UTC),
material_rate_kg_per_hour=3600.0,
gate_value=1.0,
),
MaterialSample(
timestamp=datetime(2026, 9, 3, 8, 0, 20, tzinfo=UTC),
material_rate_kg_per_hour=3600.0,
gate_value=1.0,
),
]
service = MaterialPollingService(
gateway=gateway,
state_store=state_store,
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
gate_threshold=0.5,
max_sample_gap_seconds=20.0,
)
result = service.poll_once(
now=datetime(2026, 9, 3, 8, 0, 30, tzinfo=UTC),
)
assert result is not None
gateway.get_production_runs.assert_not_called()
assert result.state.cumulative_consumption_kg == 110.0
assert result.state.integrated_running_seconds == 370.0
state_store.load.assert_called_once_with(
"machine-1",
"ORDER-1",
)
gateway.get_material_samples.assert_called_once_with(
machine_id="machine-1",
rate_variable_id="rate-variable",
gate_variable_id="gate-variable",
start=run.start,
end=datetime(2026, 9, 3, 8, 0, 30, tzinfo=UTC),
)
@pytest.fixture
def polling(tmp_path):
from production_analytics.service.material_state_store import JsonMaterialStateStore
gateway = Mock()
gateway.get_production_runs.return_value = []
gateway.get_open_production_run.return_value = EnlyzeProductionRun(
"run", "machine", None, "order", datetime(2026, 9, 3, tzinfo=UTC), None,
)
gateway.get_material_samples.return_value = []
store = JsonMaterialStateStore(tmp_path)
service = MaterialPollingService(
gateway=gateway, state_store=store, machine_id="machine",
rate_variable_id="rate", gate_variable_id="gate", gate_threshold=0.5,
max_sample_gap_seconds=20,
)
return service, gateway, store
def test_empty_response_persists_state_without_inferred_consumption(polling) -> None:
service, gateway, store = polling
now = datetime(2026, 9, 3, 1, tzinfo=UTC)
assert service.poll_once(now=now).state == MaterialIntegrationState()
saved = MaterialPollingState("run", MaterialIntegrationState(
cumulative_consumption_kg=10, integrated_running_seconds=10,
last_processed_timestamp=now, last_material_rate_kg_per_hour=3600,
last_gate_value=1, integration_active=True,
))
store.save("machine", "order", saved)
assert service.poll_once(now=now).state == saved.integration_state
assert store.load("machine", "order") == saved
def test_different_order_does_not_reuse_state(polling) -> None:
service, gateway, store = polling
store.save("machine", "other-order", MaterialPollingState(
"old", MaterialIntegrationState(cumulative_consumption_kg=100),
))
result = service.poll_once(now=datetime(2026, 9, 3, 1, tzinfo=UTC))
assert result.state == MaterialIntegrationState()
assert gateway.get_material_samples.call_args.kwargs["start"] == result.run.start
assert store.load("machine", "other-order").integration_state.cumulative_consumption_kg == 100
def test_now_before_run_start_does_not_query_or_save(polling) -> None:
service, gateway, store = polling
assert service.poll_once(now=datetime(2026, 9, 2, tzinfo=UTC)) is None
gateway.get_production_runs.assert_not_called()
gateway.get_material_samples.assert_not_called()
assert store.load("machine", "order") is None
def test_naive_now_is_rejected_before_io(polling) -> None:
service, gateway, store = polling
with pytest.raises(ValueError, match="now must be timezone-aware"):
service.poll_once(now=datetime(2026, 9, 3))
gateway.get_open_production_run.assert_not_called()
def test_new_run_empty_response_resets_baseline_and_retains_totals(polling) -> None:
service, gateway, store = polling
store.save("machine", "order", MaterialPollingState("previous", MaterialIntegrationState(
cumulative_consumption_kg=100, integrated_running_seconds=50,
last_processed_timestamp=datetime(2026, 9, 2, 23, 59, 59, tzinfo=UTC),
last_material_rate_kg_per_hour=3600, last_gate_value=1, integration_active=True,
)))
result = service.poll_once(now=datetime(2026, 9, 3, 1, tzinfo=UTC))
assert result.state == MaterialIntegrationState(
cumulative_consumption_kg=100, integrated_running_seconds=50,
)
assert store.load("machine", "order") == MaterialPollingState("run", result.state)
def test_now_before_saved_timestamp_leaves_state_unchanged(polling) -> None:
service, gateway, store = polling
saved = MaterialPollingState("run", MaterialIntegrationState(
last_processed_timestamp=datetime(2026, 9, 3, 2, tzinfo=UTC),
))
store.save("machine", "order", saved)
assert service.poll_once(now=datetime(2026, 9, 3, 1, tzinfo=UTC)) is None
gateway.get_material_samples.assert_not_called()
assert store.load("machine", "order") == saved
@pytest.mark.parametrize("order", [
"K 7-12026000815", "K 7-12025000074-K 7-12025000075",
])
@pytest.mark.parametrize("gap_seconds", [5, 86400])
def test_bootstrap_replays_exact_order_segments_and_then_resumes(
polling, order, gap_seconds,
) -> None:
from dataclasses import replace
from datetime import timedelta
service, gateway, store = polling
base = datetime(2026, 9, 3, tzinfo=UTC)
runs = [
EnlyzeProductionRun(
f"run-{i}", "machine", None, order,
base + timedelta(seconds=i * (10 + gap_seconds)),
base + timedelta(seconds=i * (10 + gap_seconds) + 10) if i < 4 else None,
)
for i in range(5)
]
current = runs[-1]
now = current.start + timedelta(seconds=10)
gateway.get_open_production_run.return_value = current
gateway.get_production_runs.return_value = [
current, runs[2], runs[0], runs[3], runs[1],
replace(runs[0], uuid="unrelated", production_order=order + "-suffix"),
replace(runs[0], uuid="component", production_order="K 7-12025000074"),
replace(runs[0], uuid="whitespace", production_order=order + " "),
replace(current, uuid="future", start=now + timedelta(days=1)),
]
def samples(**kwargs):
return [
MaterialSample(kwargs["start"], 3600, 1),
MaterialSample(kwargs["end"], 3600, 1),
]
gateway.get_material_samples.side_effect = samples
result = service.poll_once(now=now)
assert result.state.cumulative_consumption_kg == 50
assert result.state.integrated_running_seconds == 50
assert [
(call.kwargs["start"], call.kwargs["end"])
for call in gateway.get_material_samples.call_args_list
] == [(run.start, run.end or now) for run in runs]
assert store.load("machine", order) == MaterialPollingState(current.uuid, result.state)
gateway.get_production_runs.reset_mock()
gateway.get_material_samples.reset_mock()
result = service.poll_once(now=now + timedelta(seconds=10))
gateway.get_production_runs.assert_not_called()
assert gateway.get_material_samples.call_args.kwargs["start"] == now
assert result.state.cumulative_consumption_kg == 60
assert result.state.integrated_running_seconds == 60
def test_bootstrap_failure_does_not_persist_partial_totals(polling) -> None:
from dataclasses import replace
from datetime import timedelta
service, gateway, store = polling
current = gateway.get_open_production_run.return_value
gateway.get_production_runs.return_value = [
replace(current, uuid="old", start=current.start - timedelta(hours=1),
end=current.start - timedelta(minutes=30)),
current,
]
gateway.get_material_samples.side_effect = [[], ValueError("failed page")]
with pytest.raises(ValueError, match="failed page"):
service.poll_once(now=current.start + timedelta(seconds=10))
assert store.load("machine", "order") is None