410 lines
14 KiB
Python
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
|