From 4e904cac775eead00ed9b957ff5c5b81e5d17003 Mon Sep 17 00:00:00 2001 From: Brian Gallagher Date: Mon, 29 Jun 2026 21:30:02 +0000 Subject: [PATCH 01/12] Initial xcvrd CPO support - Refactor CmisManagerTask to accept a dictionary of SFP objects - Add CPO DomInfoUpdateTask and CmisManagerTask skeletons that will handle any CPO ports Signed-off-by: Brian Gallagher --- sonic-xcvrd/tests/test_cpo.py | 121 ++++++++++++++++++ sonic-xcvrd/tests/test_xcvrd.py | 130 +++++++++----------- sonic-xcvrd/xcvrd/cmis/cmis_manager_task.py | 21 +++- sonic-xcvrd/xcvrd/cpo/__init__.py | 0 sonic-xcvrd/xcvrd/cpo/cpo_manager_task.py | 22 ++++ sonic-xcvrd/xcvrd/dom/dom_mgr.py | 38 ++++-- sonic-xcvrd/xcvrd/xcvrd.py | 54 ++++---- sonic-xcvrd/xcvrd/xcvrd_utilities/common.py | 63 +++++++++- 8 files changed, 327 insertions(+), 122 deletions(-) create mode 100644 sonic-xcvrd/tests/test_cpo.py create mode 100644 sonic-xcvrd/xcvrd/cpo/__init__.py create mode 100644 sonic-xcvrd/xcvrd/cpo/cpo_manager_task.py diff --git a/sonic-xcvrd/tests/test_cpo.py b/sonic-xcvrd/tests/test_cpo.py new file mode 100644 index 000000000..b10c90231 --- /dev/null +++ b/sonic-xcvrd/tests/test_cpo.py @@ -0,0 +1,121 @@ +import sys + +if sys.version_info >= (3, 3): + from unittest.mock import MagicMock, patch +else: + from mock import MagicMock, patch + +from xcvrd.xcvrd_utilities import common + + +class TestPortDeviceResolver(object): + def test_is_cpo_port_no_chassis(self): + with patch.object(common, 'platform_chassis', None): + assert common.is_cpo_port(0) is False + + def test_is_cpo_port_true_when_cpo_present(self): + chassis = MagicMock() + chassis.get_cpo.return_value = MagicMock() + with patch.object(common, 'platform_chassis', chassis): + assert common.is_cpo_port(3) is True + chassis.get_cpo.assert_called_with(3) + + def test_is_cpo_port_false_when_not_cpo(self): + chassis = MagicMock() + chassis.get_cpo.return_value = None + with patch.object(common, 'platform_chassis', chassis): + assert common.is_cpo_port(3) is False + + def test_is_cpo_port_swallows_not_implemented(self): + chassis = MagicMock() + chassis.get_cpo.side_effect = NotImplementedError + with patch.object(common, 'platform_chassis', chassis): + assert common.is_cpo_port(3) is False + + def test_get_port_device_prefers_cpo(self): + chassis = MagicMock() + cpo = MagicMock() + chassis.get_cpo.return_value = cpo + with patch.object(common, 'platform_chassis', chassis): + assert common.get_port_device(1) is cpo + chassis.get_sfp.assert_not_called() + + def test_get_port_device_falls_back_to_sfp(self): + chassis = MagicMock() + sfp = MagicMock() + chassis.get_cpo.return_value = None + chassis.get_sfp.return_value = sfp + with patch.object(common, 'platform_chassis', chassis): + assert common.get_port_device(1) is sfp + + def test_get_port_device_none_when_unavailable(self): + with patch.object(common, 'platform_chassis', None): + assert common.get_port_device(1) is None + + +class TestObjDictAccessors(object): + def _make_port_mapping(self, physical_ports=(0, 1, 2)): + port_mapping = MagicMock() + port_mapping.physical_to_logical = {p: ['Ethernet{}'.format(p * 4)] for p in physical_ports} + return port_mapping + + def _make_obj_dict(self): + return {0: MagicMock(), 1: MagicMock(), 2: MagicMock()} + + def test_get_cpo_obj_dict(self): + objs = self._make_obj_dict() + with patch.object(common, 'is_cpo_port', side_effect=lambda p: p in (1,)), \ + patch.object(common, 'get_port_device', side_effect=lambda p: objs[p]): + cpo = common.get_cpo_obj_dict(self._make_port_mapping()) + assert set(cpo.keys()) == {1} + assert cpo[1] is objs[1] + + def test_get_pluggable_obj_dict_excludes_cpo(self): + objs = self._make_obj_dict() + with patch.object(common, 'is_cpo_port', side_effect=lambda p: p in (1,)), \ + patch.object(common, 'get_port_device', side_effect=lambda p: objs[p]): + pluggable = common.get_pluggable_obj_dict(self._make_port_mapping()) + assert set(pluggable.keys()) == {0, 2} + assert pluggable[0] is objs[0] + + def test_accessors_are_complementary(self): + objs = self._make_obj_dict() + port_mapping = self._make_port_mapping() + with patch.object(common, 'is_cpo_port', side_effect=lambda p: p in (1,)), \ + patch.object(common, 'get_port_device', side_effect=lambda p: objs[p]): + cpo = common.get_cpo_obj_dict(port_mapping) + pluggable = common.get_pluggable_obj_dict(port_mapping) + assert set(cpo) | set(pluggable) == set(objs) + assert set(cpo) & set(pluggable) == set() + + def test_all_pluggable_when_no_cpo(self): + objs = self._make_obj_dict() + port_mapping = self._make_port_mapping() + with patch.object(common, 'is_cpo_port', return_value=False), \ + patch.object(common, 'get_port_device', side_effect=lambda p: objs[p]): + assert common.get_cpo_obj_dict(port_mapping) == {} + assert set(common.get_pluggable_obj_dict(port_mapping)) == {0, 1, 2} + + def test_accessors_return_empty_without_port_mapping(self): + with patch.object(common, 'get_port_device') as mock_get_port_device: + assert common.get_cpo_obj_dict(None) == {} + assert common.get_pluggable_obj_dict(None) == {} + + port_mapping = MagicMock() + port_mapping.physical_to_logical = None + assert common.get_cpo_obj_dict(port_mapping) == {} + assert common.get_pluggable_obj_dict(port_mapping) == {} + mock_get_port_device.assert_not_called() + + def test_accessors_skip_ports_raising_exceptions(self): + objs = self._make_obj_dict() + + def mock_get_port_device(physical_port): + if physical_port == 2: + raise ValueError("Invalid port") + return objs[physical_port] + + with patch.object(common, 'is_cpo_port', return_value=False), \ + patch.object(common, 'get_port_device', side_effect=mock_get_port_device): + pluggable = common.get_pluggable_obj_dict(self._make_port_mapping()) + assert set(pluggable.keys()) == {0, 1} diff --git a/sonic-xcvrd/tests/test_xcvrd.py b/sonic-xcvrd/tests/test_xcvrd.py index a25d742ea..10e9722ea 100644 --- a/sonic-xcvrd/tests/test_xcvrd.py +++ b/sonic-xcvrd/tests/test_xcvrd.py @@ -415,7 +415,7 @@ def test_SffManagerTask_task_run_with_exception(self): def test_CmisManagerTask_task_run_with_exception(self): port_mapping = PortMapping() stop_event = threading.Event() - cmis_manager = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + cmis_manager = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) cmis_manager.wait_for_port_config_done = MagicMock(side_effect = NotImplementedError) exception_received = None trace = None @@ -434,7 +434,7 @@ def test_CmisManagerTask_task_run_with_exception(self): port_change_event = PortChangeEvent('Ethernet0', 1, 0, PortChangeEvent.PORT_ADD) port_mapping.handle_port_change_event(port_change_event) - cmis_manager = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + cmis_manager = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) cmis_manager.wait_for_port_config_done = MagicMock() #no-op cmis_manager.update_port_transceiver_status_table_sw_cmis_state = MagicMock(side_effect = NotImplementedError) exception_received = None @@ -459,16 +459,14 @@ def test_CmisManagerTask_task_run_with_exception(self): @patch('xcvrd.xcvrd_utilities.common.log_exception_traceback') @patch('xcvrd.xcvrd.XcvrTableHelper.get_status_sw_tbl') @patch('xcvrd.xcvrd.XcvrTableHelper.get_state_port_tbl') - @patch('xcvrd.xcvrd.platform_chassis') - def test_CmisManagerTask_get_xcvr_api_exception(self, mock_platform_chassis, mock_get_state_port_tbl, mock_get_status_sw_tbl, mock_log_exception_traceback): + def test_CmisManagerTask_get_xcvr_api_exception(self, mock_get_state_port_tbl, mock_get_status_sw_tbl, mock_log_exception_traceback): mock_get_status_sw_tbl = Table("STATE_DB", TRANSCEIVER_STATUS_SW_TABLE) mock_get_state_port_tbl.return_value = Table("APPL_DB", 'PORT_TABLE') mock_sfp = MagicMock() mock_sfp.get_presence.return_value = True - mock_platform_chassis.get_sfp = MagicMock(return_value=mock_sfp) port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=mock_platform_chassis) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: mock_sfp}, stop_event) task.task_stopping_event.is_set = MagicMock(side_effect=[False, False, True]) task.get_cfg_port_tbl = MagicMock() task.xcvr_table_helper = XcvrTableHelper(DEFAULT_NAMESPACE) @@ -2677,8 +2675,8 @@ def test_DaemonXcvrd_initialize_port_init_control_fields_in_port_table(self): xcvrd.initialize_port_init_control_fields_in_port_table(port_mapping) mock_state_db.set.call_count = 2 - @patch('xcvrd.xcvrd.platform_chassis') - def test_initialize_sfp_obj_dict(self, mock_platform_chassis): + @patch('xcvrd.xcvrd_utilities.common.platform_chassis') + def test_get_pluggable_obj_dict(self, mock_platform_chassis): mock_sfp_obj_1 = MagicMock() mock_sfp_obj_2 = MagicMock() def mock_get_sfp(port): @@ -2693,20 +2691,18 @@ def mock_get_sfp(port): mock_port_mapping_data = MagicMock() mock_port_mapping_data.physical_to_logical = {1: 'Ethernet0', 2: 'Ethernet1', 3: 'Ethernet2'} - # Create an instance of DaemonXcvrd - daemon_xcvrd = DaemonXcvrd(SYSLOG_IDENTIFIER) - # port_mapping is None - sfp_obj_dict = daemon_xcvrd.initialize_sfp_obj_dict(None) + sfp_obj_dict = common.get_pluggable_obj_dict(None) assert len(sfp_obj_dict) == 0 assert mock_platform_chassis.get_sfp.call_count == 0 # Mock the get_sfp method to return a MagicMock object # Call the method to test + mock_platform_chassis.get_cpo.return_value = None mock_platform_chassis.get_sfp.side_effect = mock_get_sfp - sfp_obj_dict = daemon_xcvrd.initialize_sfp_obj_dict(mock_port_mapping_data) + sfp_obj_dict = common.get_pluggable_obj_dict(mock_port_mapping_data) - # Verify the and the below also ensures that physical port 3 is not included since it is not in the port mapping + # The below also ensures that physical port 3 is not included since get_sfp raised for it assert len(sfp_obj_dict) == 2 assert 1 in sfp_obj_dict assert 2 in sfp_obj_dict @@ -3152,7 +3148,7 @@ def test_SffManagerTask_task_worker(self, mock_chassis, mock_logger): def test_CmisManagerTask_update_port_transceiver_status_table_sw_cmis_state(self): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) port_change_event = PortChangeEvent('Ethernet0', 1, 0, PortChangeEvent.PORT_SET) task.on_port_update_event(port_change_event) @@ -3170,7 +3166,7 @@ def test_CmisManagerTask_update_port_transceiver_status_table_sw_cmis_state(self def test_CmisManagerTask_handle_port_change_event(self): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) assert not task.isPortConfigDone port_change_event = PortChangeEvent('PortConfigDone', -1, 0, PortChangeEvent.PORT_SET) @@ -3210,7 +3206,7 @@ def test_CmisManagerTask_handle_port_change_event(self): def test_CmisManagerTask_get_configured_freq(self, mock_table_helper): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) cfg_port_tbl = MagicMock() cfg_port_tbl.hget = MagicMock(return_value=(True, 193100)) mock_table_helper.get_cfg_port_tbl = MagicMock(return_value=cfg_port_tbl) @@ -3222,7 +3218,7 @@ def test_CmisManagerTask_get_configured_freq(self, mock_table_helper): def test_CmisManagerTask_get_configured_tx_power_from_db(self, mock_table_helper): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) cfg_port_tbl = MagicMock() cfg_port_tbl.hget = MagicMock(return_value=(True, -10)) mock_table_helper.get_cfg_port_tbl = MagicMock(return_value=cfg_port_tbl) @@ -3231,10 +3227,9 @@ def test_CmisManagerTask_get_configured_tx_power_from_db(self, mock_table_helper assert task.get_configured_tx_power_from_db('Ethernet0') == -10 @patch('xcvrd.xcvrd.XcvrTableHelper.get_status_sw_tbl') - @patch('xcvrd.xcvrd.platform_chassis') @patch('xcvrd.xcvrd_utilities.common.is_fast_reboot_enabled', MagicMock(return_value=False)) @patch('xcvrd.xcvrd_utilities.common.get_cmis_application_desired', MagicMock(return_value=1)) - def test_CmisManagerTask_process_single_lport_invalid_host_lanes_mask(self, mock_chassis, mock_get_status_sw_tbl): + def test_CmisManagerTask_process_single_lport_invalid_host_lanes_mask(self, mock_get_status_sw_tbl): """Test process_single_lport when get_cmis_host_lanes_mask returns invalid value (<=0)""" mock_get_status_sw_tbl = Table("STATE_DB", TRANSCEIVER_STATUS_SW_TABLE) @@ -3261,13 +3256,12 @@ def test_CmisManagerTask_process_single_lport_invalid_host_lanes_mask(self, mock mock_sfp = MagicMock() mock_sfp.get_presence = MagicMock(return_value=True) mock_sfp.get_xcvr_api = MagicMock(return_value=mock_xcvr_api) - mock_chassis.get_sfp = MagicMock(return_value=mock_sfp) # Setup port mapping and task port_mapping = PortMapping() port_mapping.handle_port_change_event(PortChangeEvent('Ethernet0', 1, 0, PortChangeEvent.PORT_ADD)) stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=mock_chassis) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: mock_sfp}, stop_event) task.xcvr_table_helper = XcvrTableHelper(DEFAULT_NAMESPACE) task.xcvr_table_helper.get_status_sw_tbl.return_value = mock_get_status_sw_tbl @@ -3295,9 +3289,8 @@ def test_CmisManagerTask_process_single_lport_invalid_host_lanes_mask(self, mock assert common.get_cmis_state_from_state_db('Ethernet0', mock_get_status_sw_tbl) == CMIS_STATE_FAILED @patch('xcvrd.xcvrd.XcvrTableHelper.get_status_sw_tbl') - @patch('xcvrd.xcvrd.platform_chassis') @patch('xcvrd.xcvrd_utilities.common.is_fast_reboot_enabled', MagicMock(return_value=False)) - def test_CmisManagerTask_process_single_lport_tx_power_config_failure(self, mock_chassis, mock_get_status_sw_tbl): + def test_CmisManagerTask_process_single_lport_tx_power_config_failure(self, mock_get_status_sw_tbl): """Test process_single_lport when configure_tx_output_power fails for coherent module""" mock_get_status_sw_tbl = Table("STATE_DB", TRANSCEIVER_STATUS_SW_TABLE) @@ -3327,13 +3320,12 @@ def test_CmisManagerTask_process_single_lport_tx_power_config_failure(self, mock mock_sfp = MagicMock() mock_sfp.get_presence = MagicMock(return_value=True) mock_sfp.get_xcvr_api = MagicMock(return_value=mock_xcvr_api) - mock_chassis.get_sfp = MagicMock(return_value=mock_sfp) # Setup port mapping and task port_mapping = PortMapping() port_mapping.handle_port_change_event(PortChangeEvent('Ethernet0', 1, 0, PortChangeEvent.PORT_ADD)) stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=mock_chassis) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: mock_sfp}, stop_event) task.xcvr_table_helper = XcvrTableHelper(DEFAULT_NAMESPACE) task.xcvr_table_helper.get_status_sw_tbl.return_value = mock_get_status_sw_tbl @@ -3372,9 +3364,8 @@ def test_CmisManagerTask_process_single_lport_tx_power_config_failure(self, mock assert common.get_cmis_state_from_state_db('Ethernet0', mock_get_status_sw_tbl) == CMIS_STATE_DP_DEINIT @patch('xcvrd.xcvrd.XcvrTableHelper.get_status_sw_tbl') - @patch('xcvrd.xcvrd.platform_chassis') @patch('xcvrd.xcvrd_utilities.common.is_fast_reboot_enabled', MagicMock(return_value=False)) - def test_CmisManagerTask_dp_deinit_low_pwr_deinits_and_disables_all_lanes(self, mock_chassis, mock_get_status_sw_tbl): + def test_CmisManagerTask_dp_deinit_low_pwr_deinits_and_disables_all_lanes(self, mock_get_status_sw_tbl): """In ModuleLowPwr (first power-up), set_datapath_deinit must use max_host_lanes_mask, not the breakout subport's host_lanes_mask, because DPDeinitLanes is 0x00 after module reset. Similarly, we should disable all tx output since OutputDisableTx = 0x0 after module reset.""" @@ -3393,12 +3384,11 @@ def test_CmisManagerTask_dp_deinit_low_pwr_deinits_and_disables_all_lanes(self, mock_sfp = MagicMock() mock_sfp.get_presence = MagicMock(return_value=True) mock_sfp.get_xcvr_api = MagicMock(return_value=mock_xcvr_api) - mock_chassis.get_sfp = MagicMock(return_value=mock_sfp) port_mapping = PortMapping() port_mapping.handle_port_change_event(PortChangeEvent('Ethernet0', 1, 0, PortChangeEvent.PORT_ADD)) stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=mock_chassis) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: mock_sfp}, stop_event) task.xcvr_table_helper = XcvrTableHelper(DEFAULT_NAMESPACE) task.xcvr_table_helper.get_status_sw_tbl.return_value = mock_get_status_sw_tbl @@ -3433,7 +3423,7 @@ def test_CmisManagerTask_task_run_stop(self, mock_chassis): port_mapping = PortMapping() stop_event = threading.Event() - cmis_manager = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + cmis_manager = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) cmis_manager.wait_for_port_config_done = MagicMock() cmis_manager.start() cmis_manager.join() @@ -3460,7 +3450,7 @@ def test_CmisManagerTask_is_decommission_required(self, active_map, desired_map, mock_xcvr_api.get_application = MagicMock(side_effect=lambda lane: active_map[lane]) port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} task.get_desired_app_map = MagicMock(return_value=desired_map) assert task.is_decommission_required(mock_xcvr_api, 'Ethernet0') == expected @@ -3480,7 +3470,7 @@ def test_CmisManagerTask_is_decommission_required_uses_active_appsel_not_staged( port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} task.get_desired_app_map = MagicMock(return_value=desired_map) @@ -3503,7 +3493,7 @@ def test_CmisManagerTask_is_decommission_required_invalid_active_appsel(self): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} task.get_desired_app_map = MagicMock(return_value=desired_map) @@ -3526,7 +3516,7 @@ def test_CmisManagerTask_is_decommission_required_missing_active_appsel_lane(sel port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} task.get_desired_app_map = MagicMock(return_value=desired_map) @@ -3537,7 +3527,7 @@ def test_CmisManagerTask_get_desired_app_map(self): mock_xcvr_api = MagicMock() port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} # Mock cfg_port_tbl @@ -3580,7 +3570,7 @@ def test_CmisManagerTask_get_desired_app_map_mixed_mode(self): mock_xcvr_api = MagicMock() port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} cfg_port_tbl = MagicMock() @@ -3620,7 +3610,7 @@ def test_CmisManagerTask_get_desired_app_map_no_matching_app(self): mock_xcvr_api = MagicMock() port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} cfg_port_tbl = MagicMock() @@ -3639,7 +3629,7 @@ def test_CmisManagerTask_get_desired_app_map_uses_gearbox_lane_count(self): mock_xcvr_api = MagicMock() port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} cfg_port_tbl = MagicMock() @@ -3669,7 +3659,7 @@ def test_CmisManagerTask_get_sibling_port_configs(self): numeric fields, or get() returning not-found.""" port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} cfg_port_tbl = MagicMock() @@ -3714,7 +3704,7 @@ def test_CmisManagerTask_get_sibling_port_configs_no_cfg_port_tbl(self): """get_sibling_port_configs returns [] when cfg_port_tbl is None.""" port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} task.xcvr_table_helper.get_cfg_port_tbl = MagicMock(return_value=None) task._gearbox_lanes_dict = {} @@ -3725,7 +3715,7 @@ def test_CmisManagerTask_get_sibling_port_configs_uses_gearbox_lane_count(self): """get_sibling_port_configs uses gearbox lane count when present.""" port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} cfg_port_tbl = MagicMock() @@ -3796,7 +3786,7 @@ def get_application(lane): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) assert task.is_cmis_application_update_required(mock_xcvr_api, app_new, host_lanes_mask) == expected @@ -3865,7 +3855,7 @@ def get_host_lane_assignment_option_side_effect(app): mock_xcvr_api.get_host_lane_assignment_option = MagicMock(side_effect=get_host_lane_assignment_option_side_effect) port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) appl = common.get_cmis_application_desired(mock_xcvr_api, host_lane_count, speed) assert task.get_cmis_host_lanes_mask(mock_xcvr_api, appl, host_lane_count, subport) == expected @@ -3971,7 +3961,7 @@ def mock_get_side_effect(key): def test_CmisManagerTask_get_host_lane_count(self, gearbox_lanes_dict, lport, port_config_lanes, expected_count): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) task._gearbox_lanes_dict = gearbox_lanes_dict result = task.get_host_lane_count(lport, port_config_lanes) @@ -3981,7 +3971,7 @@ def test_CmisManagerTask_gearbox_integration_end_to_end(self): """Test end-to-end integration of gearbox line lanes with CMIS application selection""" port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) # Mock gearbox lanes dictionary - port has 4 system lanes but only 2 line lanes gearbox_lanes_dict = {"Ethernet0": 2} # 2 line lanes from gearbox @@ -4022,7 +4012,7 @@ def test_CmisManagerTask_gearbox_caching_integration(self): """Test that gearbox lanes dictionary is properly cached and used in task worker""" port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) # Mock the XcvrTableHelper to return a gearbox lanes dictionary mock_gearbox_lanes_dict = {"Ethernet0": 2, "Ethernet4": 4} @@ -4048,7 +4038,7 @@ def test_CmisManagerTask_post_port_active_apsel_to_db_error_cases(self, mock_fie port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) lport = "Ethernet0" host_lanes_mask = 0xff @@ -4110,7 +4100,7 @@ def test_CmisManagerTask_post_port_active_apsel_to_db(self): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) task.xcvr_table_helper = XcvrTableHelper(DEFAULT_NAMESPACE) task.xcvr_table_helper.get_intf_tbl = MagicMock(return_value=int_tbl) @@ -4207,7 +4197,7 @@ def test_CmisManagerTask_post_port_active_apsel_to_db(self): def test_CmisManagerTask_test_is_timer_expired(self, expired_time, current_time, expected_result): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) # Call the is_timer_expired function result = task.is_timer_expired(expired_time, current_time) @@ -4216,7 +4206,6 @@ def test_CmisManagerTask_test_is_timer_expired(self, expired_time, current_time, assert result == expected_result @patch('xcvrd.xcvrd.XcvrTableHelper.get_status_sw_tbl') - @patch('xcvrd.xcvrd.platform_chassis') @patch('xcvrd.xcvrd_utilities.common.is_fast_reboot_enabled', MagicMock(return_value=(False))) @patch('xcvrd.cmis.cmis_manager_task.PortChangeObserver', MagicMock(handle_port_update_event=MagicMock())) @patch('xcvrd.xcvrd._wrapper_get_sfp_type', MagicMock(return_value='QSFP_DD')) @@ -4225,7 +4214,7 @@ def test_CmisManagerTask_test_is_timer_expired(self, expired_time, current_time, @patch('xcvrd.xcvrd_utilities.common.is_cmis_api', MagicMock(return_value=True)) @patch('xcvrd.xcvrd_utilities.optics_si_parser.optics_si_present', MagicMock(return_value=(True))) @patch('xcvrd.xcvrd_utilities.optics_si_parser.fetch_optics_si_setting', MagicMock()) - def test_CmisManagerTask_task_worker(self, mock_chassis, mock_get_status_sw_tbl): + def test_CmisManagerTask_task_worker(self, mock_get_status_sw_tbl): mock_get_status_sw_tbl = Table("STATE_DB", TRANSCEIVER_STATUS_SW_TABLE) mock_xcvr_api = MagicMock() mock_xcvr_api.set_datapath_deinit = MagicMock(return_value=True) @@ -4399,13 +4388,11 @@ def test_CmisManagerTask_task_worker(self, mock_chassis, mock_get_status_sw_tbl) mock_sfp.get_presence = MagicMock(return_value=True) mock_sfp.get_xcvr_api = MagicMock(return_value=mock_xcvr_api) - mock_chassis.get_all_sfps = MagicMock(return_value=[mock_sfp]) - mock_chassis.get_sfp = MagicMock(return_value=mock_sfp) port_mapping = PortMapping() port_mapping.handle_port_change_event(PortChangeEvent('Ethernet0', 1, 0, PortChangeEvent.PORT_ADD)) stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=mock_chassis) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: mock_sfp}, stop_event) task.xcvr_table_helper = XcvrTableHelper(DEFAULT_NAMESPACE) task.xcvr_table_helper.get_status_sw_tbl.return_value = mock_get_status_sw_tbl task.task_stopping_event.is_set = MagicMock(side_effect=[False, False, True]) @@ -4473,7 +4460,7 @@ def test_CmisManagerTask_task_worker(self, mock_chassis, mock_get_status_sw_tbl) port_mapping = PortMapping() port_mapping.handle_port_change_event(PortChangeEvent('Ethernet1', 1, 0, PortChangeEvent.PORT_ADD)) stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=mock_chassis) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: mock_sfp}, stop_event) task.xcvr_table_helper = XcvrTableHelper(DEFAULT_NAMESPACE) task.xcvr_table_helper.get_status_sw_tbl.return_value = mock_get_status_sw_tbl task.task_stopping_event.is_set = MagicMock(side_effect=[False, False, True]) @@ -4498,14 +4485,13 @@ def test_CmisManagerTask_task_worker(self, mock_chassis, mock_get_status_sw_tbl) task.configure_laser_frequency = MagicMock(return_value=1) @patch('xcvrd.xcvrd.XcvrTableHelper.get_status_sw_tbl') - @patch('xcvrd.xcvrd.platform_chassis') @patch('xcvrd.xcvrd_utilities.common.is_fast_reboot_enabled', MagicMock(return_value=(True))) @patch('xcvrd.cmis.cmis_manager_task.PortChangeObserver', MagicMock(handle_port_update_event=MagicMock())) @patch('xcvrd.xcvrd._wrapper_get_sfp_type', MagicMock(return_value='QSFP_DD')) @patch('xcvrd.cmis.CmisManagerTask.wait_for_port_config_done', MagicMock()) @patch('xcvrd.cmis.CmisManagerTask.is_decommission_required', MagicMock(return_value=False)) @patch('xcvrd.xcvrd_utilities.common.is_cmis_api', MagicMock(return_value=True)) - def test_CmisManagerTask_task_worker_fastboot(self, mock_chassis, mock_get_status_sw_tbl): + def test_CmisManagerTask_task_worker_fastboot(self, mock_get_status_sw_tbl): mock_get_status_sw_tbl = Table("STATE_DB", TRANSCEIVER_STATUS_SW_TABLE) mock_xcvr_api = MagicMock() mock_xcvr_api.set_datapath_deinit = MagicMock(return_value=True) @@ -4599,13 +4585,11 @@ def test_CmisManagerTask_task_worker_fastboot(self, mock_chassis, mock_get_statu mock_sfp.get_presence = MagicMock(return_value=True) mock_sfp.get_xcvr_api = MagicMock(return_value=mock_xcvr_api) - mock_chassis.get_all_sfps = MagicMock(return_value=[mock_sfp]) - mock_chassis.get_sfp = MagicMock(return_value=mock_sfp) port_mapping = PortMapping() port_mapping.handle_port_change_event(PortChangeEvent('Ethernet0', 1, 0, PortChangeEvent.PORT_ADD)) stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=mock_chassis) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: mock_sfp}, stop_event) task.xcvr_table_helper = XcvrTableHelper(DEFAULT_NAMESPACE) task.xcvr_table_helper.get_status_sw_tbl.return_value = mock_get_status_sw_tbl task.task_stopping_event.is_set = MagicMock(side_effect=[False, False, True]) @@ -4638,14 +4622,13 @@ def test_CmisManagerTask_task_worker_fastboot(self, mock_chassis, mock_get_statu assert common.get_cmis_state_from_state_db('Ethernet0', task.xcvr_table_helper.get_status_sw_tbl(task.get_asic_id('Ethernet0'))) == CMIS_STATE_READY @patch('xcvrd.xcvrd.XcvrTableHelper.get_status_sw_tbl') - @patch('xcvrd.xcvrd.platform_chassis') @patch('xcvrd.xcvrd_utilities.common.is_fast_reboot_enabled', MagicMock(return_value=(False))) @patch('xcvrd.cmis.cmis_manager_task.PortChangeObserver', MagicMock(handle_port_update_event=MagicMock())) @patch('xcvrd.xcvrd._wrapper_get_sfp_type', MagicMock(return_value='QSFP_DD')) @patch('xcvrd.cmis.CmisManagerTask.wait_for_port_config_done', MagicMock()) @patch('xcvrd.cmis.CmisManagerTask.is_decommission_required', MagicMock(return_value=False)) @patch('xcvrd.xcvrd_utilities.common.is_cmis_api', MagicMock(return_value=True)) - def test_CmisManagerTask_task_worker_host_tx_ready_false_to_true(self, mock_chassis, mock_get_status_sw_tbl): + def test_CmisManagerTask_task_worker_host_tx_ready_false_to_true(self, mock_get_status_sw_tbl): mock_get_status_sw_tbl = Table("STATE_DB", TRANSCEIVER_STATUS_TABLE) mock_xcvr_api = MagicMock() dp_deinit_tx_disable_calls = [] @@ -4784,13 +4767,11 @@ def test_CmisManagerTask_task_worker_host_tx_ready_false_to_true(self, mock_chas mock_sfp.get_presence = MagicMock(return_value=True) mock_sfp.get_xcvr_api = MagicMock(return_value=mock_xcvr_api) - mock_chassis.get_all_sfps = MagicMock(return_value=[mock_sfp]) - mock_chassis.get_sfp = MagicMock(return_value=mock_sfp) port_mapping = PortMapping() port_mapping.handle_port_change_event(PortChangeEvent('Ethernet0', 1, 0, PortChangeEvent.PORT_ADD)) stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=mock_chassis) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: mock_sfp}, stop_event) task.xcvr_table_helper = XcvrTableHelper(DEFAULT_NAMESPACE) task.xcvr_table_helper.get_status_sw_tbl.return_value = mock_get_status_sw_tbl task.task_stopping_event.is_set = MagicMock(side_effect=[False, False, True]) @@ -4851,14 +4832,13 @@ def test_CmisManagerTask_task_worker_host_tx_ready_false_to_true(self, mock_chas assert task.port_dict['Ethernet0']['cmis_retries'] == 1 @patch('xcvrd.xcvrd.XcvrTableHelper.get_status_sw_tbl') - @patch('xcvrd.xcvrd.platform_chassis') @patch('xcvrd.xcvrd_utilities.common.is_fast_reboot_enabled', MagicMock(return_value=(False))) @patch('xcvrd.cmis.cmis_manager_task.PortChangeObserver', MagicMock(handle_port_update_event=MagicMock())) @patch('xcvrd.xcvrd._wrapper_get_sfp_type', MagicMock(return_value='QSFP_DD')) @patch('xcvrd.cmis.CmisManagerTask.wait_for_port_config_done', MagicMock()) @patch('xcvrd.xcvrd_utilities.common.is_cmis_api', MagicMock(return_value=True)) @patch('xcvrd.xcvrd_utilities.common.get_cmis_application_desired', MagicMock(return_value=1)) - def test_CmisManagerTask_task_worker_decommission(self, mock_chassis, mock_get_status_sw_tbl): + def test_CmisManagerTask_task_worker_decommission(self, mock_get_status_sw_tbl): mock_get_status_sw_tbl = Table("STATE_DB", TRANSCEIVER_STATUS_TABLE) mock_xcvr_api = MagicMock() mock_xcvr_api.set_datapath_deinit = MagicMock(return_value=True) @@ -4888,10 +4868,10 @@ def test_CmisManagerTask_task_worker_decommission(self, mock_chassis, mock_get_s mock_sfp = MagicMock() mock_sfp.get_presence = MagicMock(return_value=True) mock_sfp.get_xcvr_api = MagicMock(return_value=mock_xcvr_api) - mock_chassis.get_all_sfps = MagicMock(return_value=[mock_sfp]) - mock_chassis.get_sfp = MagicMock(return_value=mock_sfp) - task = CmisManagerTask(DEFAULT_NAMESPACE, PortMapping(), stop_event, platform_chassis=mock_chassis) + port_mapping = PortMapping() + + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {0: mock_sfp}, stop_event) task.xcvr_table_helper.get_status_sw_tbl.return_value = mock_get_status_sw_tbl task.is_decommission_required = MagicMock(side_effect=[True] + [False] * 20) task.get_host_tx_status = MagicMock(return_value='true') @@ -6252,7 +6232,8 @@ class MockPortMapping: @patch('xcvrd.xcvrd.DaemonXcvrd.load_platform_util', MagicMock()) @patch('xcvrd.xcvrd_utilities.port_event_helper.get_port_mapping', MagicMock(return_value=MockPortMapping)) - @patch('xcvrd.xcvrd.DaemonXcvrd.initialize_sfp_obj_dict', MagicMock()) + @patch('xcvrd.xcvrd.common.get_pluggable_obj_dict', MagicMock(return_value={})) + @patch('xcvrd.xcvrd.common.get_cpo_obj_dict', MagicMock(return_value={})) @patch('sonic_py_common.device_info.get_paths_to_platform_and_hwsku_dirs', MagicMock(return_value=('/tmp', '/tmp'))) @patch('swsscommon.swsscommon.WarmStart', MagicMock()) @patch('xcvrd.xcvrd.DaemonXcvrd.wait_for_port_config_done', MagicMock()) @@ -6299,7 +6280,8 @@ def test_DaemonXcvrd_init_deinit_fastboot_enabled(self, mock_del_port_sfp_dom_in @patch('xcvrd.xcvrd_utilities.port_event_helper.get_port_mapping', MagicMock(return_value=MockPortMapping)) @patch('sonic_py_common.device_info.get_paths_to_platform_and_hwsku_dirs', MagicMock(return_value=('/tmp', '/tmp'))) @patch('xcvrd.xcvrd.DaemonXcvrd.wait_for_port_config_done', MagicMock()) - @patch('xcvrd.xcvrd.DaemonXcvrd.initialize_sfp_obj_dict', MagicMock()) + @patch('xcvrd.xcvrd.common.get_pluggable_obj_dict', MagicMock(return_value={})) + @patch('xcvrd.xcvrd.common.get_cpo_obj_dict', MagicMock(return_value={})) @patch('subprocess.check_output', MagicMock(return_value='false')) @patch('xcvrd.xcvrd.common.is_syncd_warm_restore_complete', MagicMock(return_value=False)) @patch('xcvrd.xcvrd_utilities.common.del_port_sfp_dom_info_from_db') @@ -6447,7 +6429,7 @@ def test_CmisManagerTask_validate_frequency_and_grid(self, lport, freq, grid, ex mock_xcvr_api.get_supported_freq_config.return_value = (0x80, 0, 0, 191300, 196100) port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) result = task.validate_frequency_and_grid(mock_xcvr_api, lport, freq, grid) assert result == expected diff --git a/sonic-xcvrd/xcvrd/cmis/cmis_manager_task.py b/sonic-xcvrd/xcvrd/cmis/cmis_manager_task.py index b39992fa7..b5e191469 100644 --- a/sonic-xcvrd/xcvrd/cmis/cmis_manager_task.py +++ b/sonic-xcvrd/xcvrd/cmis/cmis_manager_task.py @@ -46,19 +46,20 @@ class CmisManagerTask(threading.Thread): CMIS_MAX_HOST_LANES = 8 CMIS_EXPIRATION_BUFFER_MS = 2 - def __init__(self, namespaces, port_mapping, main_thread_stop_event, skip_cmis_mgr=False, platform_chassis=None): + def __init__(self, namespaces, port_mapping, port_obj_dict, main_thread_stop_event, skip_cmis_mgr=False): threading.Thread.__init__(self) self.name = "CmisManagerTask" self.exc = None self.task_stopping_event = threading.Event() self.main_thread_stop_event = main_thread_stop_event self.port_dict = {k: {"asic_id": v} for k, v in port_mapping.logical_to_asic.items()} + self.port_mapping = port_mapping self.decomm_pending_dict = {} self.isPortInitDone = False self.isPortConfigDone = False self.skip_cmis_mgr = skip_cmis_mgr self.namespaces = namespaces - self.platform_chassis = platform_chassis + self.port_obj_dict = port_obj_dict self.xcvr_table_helper = XcvrTableHelper(self.namespaces) self._is_fast_reboot_enabled = None # Cache of gearbox line lanes dict, refreshed once per task_worker iteration. @@ -118,6 +119,9 @@ def on_port_update_event(self, port_change_event): if port_change_event.port_dict is None: return + if pport not in self.port_obj_dict: + return + if port_change_event.event_type == port_change_event.PORT_SET: if lport not in self.port_dict: self.port_dict[lport] = {"asic_id": port_change_event.asic_id, @@ -1272,10 +1276,13 @@ def process_single_lport(self, lport, info): if pport < 0 or speed == 0 or len(lanes) < 1 or subport < 0: return + if pport not in self.port_obj_dict: + return + host_lane_count = self.get_host_lane_count(lport, lanes) # double-check the HW presence before moving forward - sfp = self.platform_chassis.get_sfp(pport) + sfp = self.port_obj_dict[pport] if not sfp.get_presence(): self.update_port_transceiver_status_table_sw_cmis_state(lport, CMIS_STATE_REMOVED) return @@ -1351,8 +1358,8 @@ def task_worker(self): self.log_notice("Stopped") def run(self): - if self.platform_chassis is None: - self.log_notice("Platform chassis is not available, stopping...") + if not self.port_obj_dict: + self.log_notice("No SFP objects are available, stopping...") return if self.skip_cmis_mgr: @@ -1366,7 +1373,9 @@ def run(self): self.wait_for_port_config_done(namespace) for lport in self.port_dict.keys(): - self.update_port_transceiver_status_table_sw_cmis_state(lport, CMIS_STATE_UNKNOWN) + pports = self.port_mapping.get_logical_to_physical(lport) + if pports and pports[0] in self.port_obj_dict: + self.update_port_transceiver_status_table_sw_cmis_state(lport, CMIS_STATE_UNKNOWN) self.task_worker() except Exception as e: diff --git a/sonic-xcvrd/xcvrd/cpo/__init__.py b/sonic-xcvrd/xcvrd/cpo/__init__.py new file mode 100644 index 000000000..e69de29bb diff --git a/sonic-xcvrd/xcvrd/cpo/cpo_manager_task.py b/sonic-xcvrd/xcvrd/cpo/cpo_manager_task.py new file mode 100644 index 000000000..15a9d29c4 --- /dev/null +++ b/sonic-xcvrd/xcvrd/cpo/cpo_manager_task.py @@ -0,0 +1,22 @@ +#!/usr/bin/env python3 + +try: + from ..cmis.cmis_manager_task import CmisManagerTask +except ImportError as e: + raise ImportError(str(e) + " - required module not found") + + +class CpoManagerTask(CmisManagerTask): + def __init__(self, namespaces, port_mapping, port_obj_dict, main_thread_stop_event, skip_cpo_mgr=False): + super().__init__(namespaces, port_mapping, port_obj_dict, main_thread_stop_event, + skip_cmis_mgr=skip_cpo_mgr) + self.name = "CpoManagerTask" + + def log_debug(self, message): + super().log_debug("CPO: {}".format(message)) + + def log_notice(self, message): + super().log_notice("CPO: {}".format(message)) + + def log_error(self, message): + super().log_error("CPO: {}".format(message)) diff --git a/sonic-xcvrd/xcvrd/dom/dom_mgr.py b/sonic-xcvrd/xcvrd/dom/dom_mgr.py index 2ca37dcf7..063a348cc 100644 --- a/sonic-xcvrd/xcvrd/dom/dom_mgr.py +++ b/sonic-xcvrd/xcvrd/dom/dom_mgr.py @@ -40,7 +40,7 @@ class DomInfoUpdateBase(threading.Thread): name = '' - def __init__(self, namespaces, port_mapping, sfp_obj_dict, main_thread_stop_event): + def __init__(self, namespaces, port_mapping, port_obj_dict, main_thread_stop_event): threading.Thread.__init__(self) self.exc = None self.task_stopping_event = threading.Event() @@ -48,7 +48,7 @@ def __init__(self, namespaces, port_mapping, sfp_obj_dict, main_thread_stop_even self.helper_logger = syslogger.SysLogger(SYSLOG_IDENTIFIER_DOMINFOUPDATETASK, enable_runtime_config=True) self.port_mapping = copy.deepcopy(port_mapping) self.namespaces = namespaces - self.sfp_obj_dict = sfp_obj_dict + self.port_obj_dict = port_obj_dict def log_debug(self, message): self.helper_logger.log_debug("{}".format(message)) @@ -147,17 +147,17 @@ class DomInfoUpdateTask(DomInfoUpdateBase): {'APPL_DB': 'PORT_TABLE', 'FILTER': ['flap_count']}, ] - def __init__(self, namespaces, port_mapping, sfp_obj_dict, main_thread_stop_event, skip_cmis_mgr, dom_update_interval=None): - super().__init__(namespaces, port_mapping, sfp_obj_dict, main_thread_stop_event) + def __init__(self, namespaces, port_mapping, port_obj_dict, main_thread_stop_event, skip_cmis_mgr, dom_update_interval=None): + super().__init__(namespaces, port_mapping, port_obj_dict, main_thread_stop_event) self.skip_cmis_mgr = skip_cmis_mgr self.link_change_affected_ports = {} self.xcvr_table_helper = XcvrTableHelper(self.namespaces) - self.xcvrd_utils = XCVRDUtils(self.sfp_obj_dict, self.helper_logger) - self.dom_db_utils = DOMDBUtils(self.sfp_obj_dict, self.port_mapping, self.xcvr_table_helper, self.task_stopping_event, self.helper_logger) + self.xcvrd_utils = XCVRDUtils(self.port_obj_dict, self.helper_logger) + self.dom_db_utils = DOMDBUtils(self.port_obj_dict, self.port_mapping, self.xcvr_table_helper, self.task_stopping_event, self.helper_logger) self.db_utils = self.dom_db_utils - self.vdm_utils = VDMUtils(self.sfp_obj_dict, self.helper_logger) - self.vdm_db_utils = VDMDBUtils(self.sfp_obj_dict, self.port_mapping, self.xcvr_table_helper, self.task_stopping_event, self.helper_logger) - self.status_db_utils = StatusDBUtils(self.sfp_obj_dict, self.port_mapping, self.xcvr_table_helper, self.task_stopping_event, self.helper_logger) + self.vdm_utils = VDMUtils(self.port_obj_dict, self.helper_logger) + self.vdm_db_utils = VDMDBUtils(self.port_obj_dict, self.port_mapping, self.xcvr_table_helper, self.task_stopping_event, self.helper_logger) + self.status_db_utils = StatusDBUtils(self.port_obj_dict, self.port_mapping, self.xcvr_table_helper, self.task_stopping_event, self.helper_logger) self.dom_update_interval = self.DEFAULT_DOM_INFO_UPDATE_PERIOD_SECS if dom_update_interval is not None: if dom_update_interval < 0: @@ -329,6 +329,9 @@ def task_worker(self): self.log_notice("Stop event generated during DOM monitoring loop") break + if physical_port not in self.port_obj_dict: + continue + # Get the first logical port name since it corresponds to the first subport # of the breakout group logical_port_name = logical_ports[0] @@ -443,6 +446,9 @@ def update_port_db_diagnostics_on_link_change(self, physical_port): if self.task_stopping_event.is_set(): return + if physical_port not in self.port_obj_dict: + return + logical_port_list = self.port_mapping.get_physical_to_logical(physical_port) if logical_port_list is None: self.log_warning("Update DB diagnostics during link change: Unknown physical port index {}".format(physical_port)) @@ -523,14 +529,19 @@ def on_remove_logical_port(self, port_change_event): self.xcvr_table_helper.get_firmware_info_tbl(port_change_event.asic_id) ]) + +class CpoDomInfoUpdateTask(DomInfoUpdateTask): + name = "CpoDomInfoUpdateTask" + + class DomThermalInfoUpdateTask(DomInfoUpdateBase): name = 'DomThermalInfoUpdateTask' - def __init__(self, namespaces, port_mapping, sfp_obj_dict, main_thread_stop_event, poll_interval): - super().__init__(namespaces, port_mapping, sfp_obj_dict, main_thread_stop_event) + def __init__(self, namespaces, port_mapping, port_obj_dict, main_thread_stop_event, poll_interval): + super().__init__(namespaces, port_mapping, port_obj_dict, main_thread_stop_event) self.poll_interval = poll_interval self.xcvr_table_helper = XcvrTableHelper(self.namespaces) - self.dom_db_utils = DOMDBUtils(self.sfp_obj_dict, self.port_mapping, self.xcvr_table_helper, self.task_stopping_event, self.helper_logger) + self.dom_db_utils = DOMDBUtils(self.port_obj_dict, self.port_mapping, self.xcvr_table_helper, self.task_stopping_event, self.helper_logger) def task_worker(self): self.log_notice("Start DOM thermal monitoring loop") @@ -551,6 +562,9 @@ def task_worker(self): continue for physical_port, logical_ports in self.port_mapping.physical_to_logical.items(): + if physical_port not in self.port_obj_dict: + continue + # Get the first logical port name since it corresponds to the first subport # of the breakout group logical_port_name = logical_ports[0] diff --git a/sonic-xcvrd/xcvrd/xcvrd.py b/sonic-xcvrd/xcvrd/xcvrd.py index ee1545334..f76510222 100644 --- a/sonic-xcvrd/xcvrd/xcvrd.py +++ b/sonic-xcvrd/xcvrd/xcvrd.py @@ -28,8 +28,9 @@ from .xcvrd_utilities import sfp_status_helper from .sff_mgr import SffManagerTask - from .dom.dom_mgr import DomThermalInfoUpdateTask, DomInfoUpdateTask + from .dom.dom_mgr import DomThermalInfoUpdateTask, DomInfoUpdateTask, CpoDomInfoUpdateTask from .cmis.cmis_manager_task import CmisManagerTask + from .cpo.cpo_manager_task import CpoManagerTask from .xcvrd_utilities.xcvr_table_helper import * from .xcvrd_utilities import port_event_helper from .xcvrd_utilities.port_event_helper import PortChangeObserver @@ -875,17 +876,19 @@ def update_log_level(self): class DaemonXcvrd(daemon_base.DaemonBase): - def __init__(self, log_identifier, skip_cmis_mgr=False, enable_sff_mgr=False, dom_temperature_poll_interval=None, dom_update_interval=None): + def __init__(self, log_identifier, skip_cmis_mgr=False, enable_sff_mgr=False, dom_temperature_poll_interval=None, dom_update_interval=None, skip_cpo_mgr=False): super(DaemonXcvrd, self).__init__(log_identifier, enable_runtime_log_config=True) self.stop_event = threading.Event() self.sfp_error_event = threading.Event() self.skip_cmis_mgr = skip_cmis_mgr + self.skip_cpo_mgr = skip_cpo_mgr self.enable_sff_mgr = enable_sff_mgr self.dom_temperature_poll_interval = dom_temperature_poll_interval self.dom_update_interval = dom_update_interval self.namespaces = [''] self.threads = [] self.sfp_obj_dict = {} + self.cpo_obj_dict = {} def update_loggers_log_level(self): """ @@ -959,30 +962,6 @@ def initialize_port_init_control_fields_in_port_table(self, port_mapping_data): self.log_notice("XCVRD INIT: Port init control fields initialized in STATE_DB PORT_TABLE") - def initialize_sfp_obj_dict(self, port_mapping_data): - """ - Create a dictionary mapping physical ports to their corresponding SFP objects. - - Args: - port_mapping_data (PortMapping): The port mapping data. - - Returns: - Dict[int, Sfp]: A dictionary mapping physical ports to SFP objects. - """ - if port_mapping_data is None or port_mapping_data.physical_to_logical is None: - self.log_error("SFP OBJ INIT: Failed to get port mapping data") - return {} - - physical_port_list = port_mapping_data.physical_to_logical.keys() - sfp_obj_dict = {} - for physical_port in physical_port_list: - try: - sfp_obj_dict[physical_port] = platform_chassis.get_sfp(physical_port) - except Exception as e: - self.log_error(f"SFP OBJ INIT: Failed to get SFP object for port {physical_port} due to {repr(e)}") - - return sfp_obj_dict - def remove_stale_transceiver_info(self, port_mapping_data): """ Remove stale entries from the TRANSCEIVER_INFO table for ports where the transceiver is no longer present. @@ -1068,7 +1047,8 @@ def init(self): port_mapping_data = port_event_helper.get_port_mapping(self.namespaces) self.initialize_port_init_control_fields_in_port_table(port_mapping_data) - self.sfp_obj_dict = self.initialize_sfp_obj_dict(port_mapping_data) + self.sfp_obj_dict = common.get_pluggable_obj_dict(port_mapping_data) + self.cpo_obj_dict = common.get_cpo_obj_dict(port_mapping_data) # Remove the TRANSCEIVER_INFO table if the transceiver is absent. # This ensures stale entries are cleaned up when a transceiver is removed while xcvrd is not running. @@ -1157,15 +1137,29 @@ def run(self): # Start the CMIS manager cmis_manager = None if not self.skip_cmis_mgr: - cmis_manager = CmisManagerTask(self.namespaces, port_mapping_data, self.stop_event, skip_cmis_mgr=self.skip_cmis_mgr, platform_chassis=platform_chassis) + cmis_manager = CmisManagerTask(self.namespaces, port_mapping_data, self.sfp_obj_dict, self.stop_event, skip_cmis_mgr=self.skip_cmis_mgr) cmis_manager.start() self.threads.append(cmis_manager) + # Start the CPO manager + cpo_manager = None + if self.cpo_obj_dict and not self.skip_cpo_mgr: + cpo_manager = CpoManagerTask(self.namespaces, port_mapping_data, self.cpo_obj_dict, self.stop_event, skip_cpo_mgr=self.skip_cpo_mgr) + cpo_manager.start() + self.threads.append(cpo_manager) + # Start the dom sensor info update thread dom_info_update = DomInfoUpdateTask(self.namespaces, port_mapping_data, self.sfp_obj_dict, self.stop_event, self.skip_cmis_mgr, self.dom_update_interval) dom_info_update.start() self.threads.append(dom_info_update) + # Start the CPO dom sensor info update thread + cpo_dom_info_update = None + if self.cpo_obj_dict: + cpo_dom_info_update = CpoDomInfoUpdateTask(self.namespaces, port_mapping_data, self.cpo_obj_dict, self.stop_event, False, self.dom_update_interval) + cpo_dom_info_update.start() + self.threads.append(cpo_dom_info_update) + # Start the dom thermal sensor info update thread dom_thermal_info_update = None if self.dom_temperature_poll_interval is not None: @@ -1245,13 +1239,15 @@ def run(self): def main(): parser = argparse.ArgumentParser() parser.add_argument('--skip_cmis_mgr', action='store_true') + parser.add_argument('--skip_cpo_mgr', action='store_true') parser.add_argument('--enable_sff_mgr', action='store_true') parser.add_argument('--dom_temperature_poll_interval', default=None, type=int) parser.add_argument('--dom_update_interval', default=None, type=int) args = parser.parse_args() xcvrd = DaemonXcvrd(SYSLOG_IDENTIFIER, args.skip_cmis_mgr, args.enable_sff_mgr, - args.dom_temperature_poll_interval, args.dom_update_interval) + args.dom_temperature_poll_interval, args.dom_update_interval, + args.skip_cpo_mgr) xcvrd.run() diff --git a/sonic-xcvrd/xcvrd/xcvrd_utilities/common.py b/sonic-xcvrd/xcvrd/xcvrd_utilities/common.py index 7a5c84d7a..83c44425e 100644 --- a/sonic-xcvrd/xcvrd/xcvrd_utilities/common.py +++ b/sonic-xcvrd/xcvrd/xcvrd_utilities/common.py @@ -112,6 +112,65 @@ def update_port_transceiver_status_table_sw(logical_port_name, status_sw_tbl, st fvs = swsscommon.FieldValuePairs([('status', status), ('error', error_descriptions)]) status_sw_tbl.set(logical_port_name, fvs) +def get_port_device(physical_port): + if platform_chassis is None: + return None + try: + cpo = platform_chassis.get_cpo(physical_port) + if cpo is not None: + return cpo + except (NotImplementedError, AttributeError, IndexError): + pass + try: + return platform_chassis.get_sfp(physical_port) + except (NotImplementedError, AttributeError, IndexError): + return None + +def is_cpo_port(physical_port): + if platform_chassis is None: + return False + try: + return platform_chassis.get_cpo(physical_port) is not None + except (NotImplementedError, AttributeError, IndexError): + return False + +def is_pluggable_port(physical_port): + return not is_cpo_port(physical_port) + +def _get_port_obj_dict(port_mapping_data, port_filter): + """ + Create a dictionary mapping physical ports to their corresponding device objects, + restricted to the ports accepted by port_filter. + + Args: + port_mapping_data (PortMapping): The port mapping data. + port_filter (Callable[[int], bool]): Predicate selecting the physical ports to include. + + Returns: + Dict[int, object]: A dictionary mapping physical ports to device objects. + """ + if port_mapping_data is None or port_mapping_data.physical_to_logical is None: + helper_logger.log_error("PORT OBJ INIT: Failed to get port mapping data") + return {} + + obj_dict = {} + for physical_port in port_mapping_data.physical_to_logical.keys(): + try: + if port_filter(physical_port): + obj_dict[physical_port] = get_port_device(physical_port) + except Exception as e: + helper_logger.log_error(f"PORT OBJ INIT: Failed to get device object for port {physical_port} due to {repr(e)}") + + return obj_dict + +def get_cpo_obj_dict(port_mapping_data): + """Create a dictionary mapping physical ports to their corresponding CPO objects.""" + return _get_port_obj_dict(port_mapping_data, is_cpo_port) + +def get_pluggable_obj_dict(port_mapping_data): + """Create a dictionary mapping physical ports to their corresponding SFP objects.""" + return _get_port_obj_dict(port_mapping_data, is_pluggable_port) + def is_copper(physical_port): """Check if the transceiver on the given physical port is copper""" if platform_chassis: @@ -125,7 +184,9 @@ def _wrapper_get_presence(physical_port): """Wrapper function to get SFP presence status""" if platform_chassis is not None: try: - return platform_chassis.get_sfp(physical_port).get_presence() + device = get_port_device(physical_port) + if device is not None: + return device.get_presence() except NotImplementedError: pass if platform_sfputil is not None: From df83f8b4346bd5deade70aa1fc0fd4fbb583839e Mon Sep 17 00:00:00 2001 From: Brian Gallagher Date: Mon, 27 Jul 2026 23:28:21 +0000 Subject: [PATCH 02/12] Do not create CmisManagerTask if there are no pluggable ports Signed-off-by: Brian Gallagher --- sonic-xcvrd/xcvrd/xcvrd.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sonic-xcvrd/xcvrd/xcvrd.py b/sonic-xcvrd/xcvrd/xcvrd.py index f76510222..0c6915ee6 100644 --- a/sonic-xcvrd/xcvrd/xcvrd.py +++ b/sonic-xcvrd/xcvrd/xcvrd.py @@ -1136,7 +1136,7 @@ def run(self): # Start the CMIS manager cmis_manager = None - if not self.skip_cmis_mgr: + if self.sfp_obj_dict and not self.skip_cmis_mgr: cmis_manager = CmisManagerTask(self.namespaces, port_mapping_data, self.sfp_obj_dict, self.stop_event, skip_cmis_mgr=self.skip_cmis_mgr) cmis_manager.start() self.threads.append(cmis_manager) From 5f2d9ae4654204a0624c0cd91538aefd5ca73fa5 Mon Sep 17 00:00:00 2001 From: Brian Gallagher Date: Tue, 28 Jul 2026 00:10:10 +0000 Subject: [PATCH 03/12] Add a CpoStateUpdateTask Signed-off-by: Brian Gallagher --- sonic-xcvrd/xcvrd/xcvrd.py | 29 +++++++++++++++++++++++++++++ 1 file changed, 29 insertions(+) diff --git a/sonic-xcvrd/xcvrd/xcvrd.py b/sonic-xcvrd/xcvrd/xcvrd.py index 0c6915ee6..eec1711eb 100644 --- a/sonic-xcvrd/xcvrd/xcvrd.py +++ b/sonic-xcvrd/xcvrd/xcvrd.py @@ -870,6 +870,12 @@ def update_log_level(self): return self.logger.update_log_level() +class CpoStateUpdateTask(SfpStateUpdateTask): + def __init__(self, namespaces, port_mapping, port_obj_dict, main_thread_stop_event, sfp_error_event): + super().__init__(namespaces, port_mapping, port_obj_dict, main_thread_stop_event, sfp_error_event) + self.name = "CpoStateUpdateTask" + + # # Daemon ======================================================================= # @@ -1173,6 +1179,13 @@ def run(self): sfp_state_update.start() self.threads.append(sfp_state_update) + # Start the CPO state info update thread + cpo_state_update = None + if self.cpo_obj_dict: + cpo_state_update = CpoStateUpdateTask(self.namespaces, port_mapping_data, self.cpo_obj_dict, self.stop_event, self.sfp_error_event) + cpo_state_update.start() + self.threads.append(cpo_state_update) + # Start main loop self.log_notice("Start daemon main loop with thread count {}".format(len(self.threads))) for thread in self.threads: @@ -1206,10 +1219,20 @@ def run(self): if cmis_manager.is_alive(): cmis_manager.join() + # Stop the CPO manager + if cpo_manager is not None: + if cpo_manager.is_alive(): + cpo_manager.join() + # Stop the dom sensor info update thread if dom_info_update.is_alive(): dom_info_update.join() + # Stop the CPO dom sensor info update thread + if cpo_dom_info_update is not None: + if cpo_dom_info_update.is_alive(): + cpo_dom_info_update.join() + # Stop the dom thermal sensor info update thread if dom_thermal_info_update is not None: if dom_thermal_info_update.is_alive(): @@ -1220,6 +1243,12 @@ def run(self): sfp_state_update.raise_exception() sfp_state_update.join() + # Stop the CPO state info update thread + if cpo_state_update is not None: + if cpo_state_update.is_alive(): + cpo_state_update.raise_exception() + cpo_state_update.join() + # Start daemon deinitialization sequence self.deinit() From d50896b88747b9a6f4de0284945505ca67232819 Mon Sep 17 00:00:00 2001 From: Brian Gallagher Date: Tue, 28 Jul 2026 00:22:39 +0000 Subject: [PATCH 04/12] Rename sfp_obj_dict to port_obj_dict in SfpStateUpdateTask Signed-off-by: Brian Gallagher --- sonic-xcvrd/xcvrd/xcvrd.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/sonic-xcvrd/xcvrd/xcvrd.py b/sonic-xcvrd/xcvrd/xcvrd.py index eec1711eb..5b115145c 100644 --- a/sonic-xcvrd/xcvrd/xcvrd.py +++ b/sonic-xcvrd/xcvrd/xcvrd.py @@ -259,7 +259,7 @@ def waiting_time_compensation_with_sleep(time_start, time_to_wait): class SfpStateUpdateTask(threading.Thread): RETRY_EEPROM_READING_INTERVAL = 60 - def __init__(self, namespaces, port_mapping, sfp_obj_dict, main_thread_stop_event, sfp_error_event): + def __init__(self, namespaces, port_mapping, port_obj_dict, main_thread_stop_event, sfp_error_event): threading.Thread.__init__(self) self.name = "SfpStateUpdateTask" self.exc = None @@ -276,11 +276,11 @@ def __init__(self, namespaces, port_mapping, sfp_obj_dict, main_thread_stop_even self.sfp_error_dict = {} self.sfp_insert_events = {} self.namespaces = namespaces - self.sfp_obj_dict = sfp_obj_dict + self.port_obj_dict = port_obj_dict self.logger = syslogger.SysLogger(SYSLOG_IDENTIFIER_SFPSTATEUPDATETASK, enable_runtime_config=True) self.xcvr_table_helper = XcvrTableHelper(self.namespaces) - self.dom_db_utils = DOMDBUtils(sfp_obj_dict, self.port_mapping, self.xcvr_table_helper, self.task_stopping_event, self.logger) - self.vdm_db_utils = VDMDBUtils(sfp_obj_dict, self.port_mapping, self.xcvr_table_helper, self.task_stopping_event, self.logger) + self.dom_db_utils = DOMDBUtils(port_obj_dict, self.port_mapping, self.xcvr_table_helper, self.task_stopping_event, self.logger) + self.vdm_db_utils = VDMDBUtils(port_obj_dict, self.port_mapping, self.xcvr_table_helper, self.task_stopping_event, self.logger) def _mapping_event_from_change_event(self, status, port_dict): """ From 2a6c1ece69f2f3a3ab590618dc3af4641696bd79 Mon Sep 17 00:00:00 2001 From: Brian Gallagher Date: Tue, 28 Jul 2026 18:35:46 +0000 Subject: [PATCH 05/12] Fix tests Signed-off-by: Brian Gallagher --- sonic-xcvrd/tests/test_xcvrd.py | 20 ++++++++++++-------- 1 file changed, 12 insertions(+), 8 deletions(-) diff --git a/sonic-xcvrd/tests/test_xcvrd.py b/sonic-xcvrd/tests/test_xcvrd.py index 10e9722ea..8f4052cc0 100644 --- a/sonic-xcvrd/tests/test_xcvrd.py +++ b/sonic-xcvrd/tests/test_xcvrd.py @@ -415,7 +415,7 @@ def test_SffManagerTask_task_run_with_exception(self): def test_CmisManagerTask_task_run_with_exception(self): port_mapping = PortMapping() stop_event = threading.Event() - cmis_manager = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + cmis_manager = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) cmis_manager.wait_for_port_config_done = MagicMock(side_effect = NotImplementedError) exception_received = None trace = None @@ -434,7 +434,7 @@ def test_CmisManagerTask_task_run_with_exception(self): port_change_event = PortChangeEvent('Ethernet0', 1, 0, PortChangeEvent.PORT_ADD) port_mapping.handle_port_change_event(port_change_event) - cmis_manager = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + cmis_manager = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) cmis_manager.wait_for_port_config_done = MagicMock() #no-op cmis_manager.update_port_transceiver_status_table_sw_cmis_state = MagicMock(side_effect = NotImplementedError) exception_received = None @@ -583,6 +583,9 @@ def test_DaemonXcvrd_run_with_exception(self, mock_task_join_sff, mock_task_join xcvrd.dom_temperature_poll_interval = 10 xcvrd.load_feature_flags = MagicMock() xcvrd.stop_event.wait = MagicMock() + # init() is mocked out, so populate the pluggable port dict directly. + # CmisManagerTask is only created when there is at least one pluggable port. + xcvrd.sfp_obj_dict = {1: MagicMock()} xcvrd.run() assert len(xcvrd.threads) == 5 @@ -3166,7 +3169,7 @@ def test_CmisManagerTask_update_port_transceiver_status_table_sw_cmis_state(self def test_CmisManagerTask_handle_port_change_event(self): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) assert not task.isPortConfigDone port_change_event = PortChangeEvent('PortConfigDone', -1, 0, PortChangeEvent.PORT_SET) @@ -5087,7 +5090,7 @@ def test_DomInfoUpdateTask_task_worker(self, mock_post_pm_info, mock_sub_table.return_value = mock_selectable port_mapping = PortMapping() - mock_sfp_obj_dict = MagicMock() + mock_sfp_obj_dict = {'1': MagicMock()} stop_event = threading.Event() mock_cmis_manager = MagicMock() task = DomInfoUpdateTask(DEFAULT_NAMESPACE, port_mapping, mock_sfp_obj_dict, stop_event, mock_cmis_manager, 0) @@ -5145,7 +5148,7 @@ def test_DomInfoUpdateTask_task_worker(self, mock_post_pm_info, @patch('xcvrd.dom.dom_mgr.DomInfoUpdateTask.post_port_pm_info_to_db') def test_DomInfoUpdateTask_task_worker_vdm_failure(self, mock_post_pm_info): port_mapping = PortMapping() - mock_sfp_obj_dict = MagicMock() + mock_sfp_obj_dict = {'1': MagicMock()} stop_event = threading.Event() mock_cmis_manager = MagicMock() task = DomInfoUpdateTask(DEFAULT_NAMESPACE, port_mapping, mock_sfp_obj_dict, stop_event, mock_cmis_manager, 0) @@ -5296,7 +5299,7 @@ def mock_log_notice(message): def test_DomInfoUpdateTask_task_worker_vdm_freeze_conditions(self, mock_post_pm_info): """Test various need_freeze condition combinations""" port_mapping = PortMapping() - mock_sfp_obj_dict = MagicMock() + mock_sfp_obj_dict = {'1': MagicMock()} stop_event = threading.Event() mock_cmis_manager = MagicMock() @@ -5444,7 +5447,7 @@ def test_update_port_db_diagnostics_on_link_change( expected_logs, ): port_mapping = PortMapping() - mock_sfp_obj_dict = MagicMock() + mock_sfp_obj_dict = {physical_port: MagicMock()} stop_event = threading.Event() mock_cmis_manager = MagicMock() task = DomInfoUpdateTask(DEFAULT_NAMESPACE, port_mapping, mock_sfp_obj_dict, stop_event, mock_cmis_manager) @@ -5829,6 +5832,7 @@ def test_wrapper_get_presence(self, mock_sfputil, mock_chassis): mock_object = MagicMock() mock_object.get_presence = MagicMock(return_value=True) mock_chassis.get_sfp = MagicMock(return_value=mock_object) + mock_chassis.get_cpo = MagicMock(return_value=None) from xcvrd.xcvrd_utilities.common import _wrapper_get_presence assert _wrapper_get_presence(1) @@ -6705,7 +6709,7 @@ def test_DomInfoUpdateTask_scheduling_uses_loop_start_time(self, mock_post_pm_in - If wrong: update doesn't trigger yet (61 < 65), which we can detect """ port_mapping = PortMapping() - mock_sfp_obj_dict = MagicMock() + mock_sfp_obj_dict = {1: MagicMock()} stop_event = threading.Event() mock_cmis_manager = MagicMock() task = DomInfoUpdateTask(DEFAULT_NAMESPACE, port_mapping, mock_sfp_obj_dict, stop_event, mock_cmis_manager) From e9c88103399a7ae685f603291f953cbb6a9ac850 Mon Sep 17 00:00:00 2001 From: Brian Gallagher Date: Tue, 28 Jul 2026 21:36:39 +0000 Subject: [PATCH 06/12] Meet test coverage threshold Signed-off-by: Brian Gallagher --- sonic-xcvrd/tests/test_xcvrd.py | 27 +++++++++++++++++++++++ sonic-xcvrd/xcvrd/cpo/cpo_manager_task.py | 9 -------- 2 files changed, 27 insertions(+), 9 deletions(-) diff --git a/sonic-xcvrd/tests/test_xcvrd.py b/sonic-xcvrd/tests/test_xcvrd.py index 8f4052cc0..d3d7254f8 100644 --- a/sonic-xcvrd/tests/test_xcvrd.py +++ b/sonic-xcvrd/tests/test_xcvrd.py @@ -2777,6 +2777,33 @@ def test_DaemonXcvrd_run(self, mock_task_stop1, mock_task_stop2, mock_task_run1, assert mock_deinit.call_count == 1 assert mock_init.call_count == 1 + @patch('xcvrd.xcvrd.SfpStateUpdateTask.is_alive', MagicMock(return_value=False)) + @patch('xcvrd.xcvrd.SfpStateUpdateTask.join', MagicMock()) + @patch('xcvrd.xcvrd.SfpStateUpdateTask.start', MagicMock()) + @patch('xcvrd.xcvrd.DomInfoUpdateTask.is_alive', MagicMock(return_value=False)) + @patch('xcvrd.xcvrd.DomInfoUpdateTask.join', MagicMock()) + @patch('xcvrd.xcvrd.DomInfoUpdateTask.start', MagicMock()) + @patch('xcvrd.cmis.CmisManagerTask.is_alive', MagicMock(return_value=False)) + @patch('xcvrd.cmis.CmisManagerTask.join', MagicMock()) + @patch('xcvrd.cmis.CmisManagerTask.start', MagicMock()) + @patch('xcvrd.xcvrd.DaemonXcvrd.deinit', MagicMock()) + @patch('xcvrd.xcvrd.DaemonXcvrd.init') + def test_DaemonXcvrd_run_creates_cpo_and_pluggable_tasks(self, mock_init): + mock_init.return_value = PortMapping() + xcvrd = DaemonXcvrd(SYSLOG_IDENTIFIER) + xcvrd.load_feature_flags = MagicMock() + xcvrd.stop_event.wait = MagicMock() + # init() is mocked out, so populate the port dicts it would normally build + xcvrd.sfp_obj_dict = {1: MagicMock()} + xcvrd.cpo_obj_dict = {2: MagicMock()} + xcvrd.run() + + # A mix of pluggable and CPO ports must yield both sets of tasks + names = [thread.name for thread in xcvrd.threads] + for expected in ('CmisManagerTask', 'DomInfoUpdateTask', 'SfpStateUpdateTask', + 'CpoManagerTask', 'CpoDomInfoUpdateTask', 'CpoStateUpdateTask'): + assert expected in names + def test_SffManagerTask_handle_port_change_event(self): stop_event = threading.Event() task = SffManagerTask(DEFAULT_NAMESPACE, stop_event, MagicMock(), helper_logger) diff --git a/sonic-xcvrd/xcvrd/cpo/cpo_manager_task.py b/sonic-xcvrd/xcvrd/cpo/cpo_manager_task.py index 15a9d29c4..7ef78aff6 100644 --- a/sonic-xcvrd/xcvrd/cpo/cpo_manager_task.py +++ b/sonic-xcvrd/xcvrd/cpo/cpo_manager_task.py @@ -11,12 +11,3 @@ def __init__(self, namespaces, port_mapping, port_obj_dict, main_thread_stop_eve super().__init__(namespaces, port_mapping, port_obj_dict, main_thread_stop_event, skip_cmis_mgr=skip_cpo_mgr) self.name = "CpoManagerTask" - - def log_debug(self, message): - super().log_debug("CPO: {}".format(message)) - - def log_notice(self, message): - super().log_notice("CPO: {}".format(message)) - - def log_error(self, message): - super().log_error("CPO: {}".format(message)) From bb38069c0bf60c5463e23137c23615835187aead Mon Sep 17 00:00:00 2001 From: Brian Gallagher Date: Wed, 29 Jul 2026 05:21:20 +0000 Subject: [PATCH 07/12] Remove some unnecessary test code Signed-off-by: Brian Gallagher --- sonic-xcvrd/tests/test_cpo.py | 29 ++++----------------- sonic-xcvrd/xcvrd/xcvrd_utilities/common.py | 5 +++- 2 files changed, 9 insertions(+), 25 deletions(-) diff --git a/sonic-xcvrd/tests/test_cpo.py b/sonic-xcvrd/tests/test_cpo.py index b10c90231..a55f9a52e 100644 --- a/sonic-xcvrd/tests/test_cpo.py +++ b/sonic-xcvrd/tests/test_cpo.py @@ -1,14 +1,9 @@ -import sys - -if sys.version_info >= (3, 3): - from unittest.mock import MagicMock, patch -else: - from mock import MagicMock, patch +from unittest.mock import MagicMock, patch from xcvrd.xcvrd_utilities import common -class TestPortDeviceResolver(object): +class TestPortDeviceResolver: def test_is_cpo_port_no_chassis(self): with patch.object(common, 'platform_chassis', None): assert common.is_cpo_port(0) is False @@ -53,7 +48,7 @@ def test_get_port_device_none_when_unavailable(self): assert common.get_port_device(1) is None -class TestObjDictAccessors(object): +class TestObjDictAccessors: def _make_port_mapping(self, physical_ports=(0, 1, 2)): port_mapping = MagicMock() port_mapping.physical_to_logical = {p: ['Ethernet{}'.format(p * 4)] for p in physical_ports} @@ -62,22 +57,6 @@ def _make_port_mapping(self, physical_ports=(0, 1, 2)): def _make_obj_dict(self): return {0: MagicMock(), 1: MagicMock(), 2: MagicMock()} - def test_get_cpo_obj_dict(self): - objs = self._make_obj_dict() - with patch.object(common, 'is_cpo_port', side_effect=lambda p: p in (1,)), \ - patch.object(common, 'get_port_device', side_effect=lambda p: objs[p]): - cpo = common.get_cpo_obj_dict(self._make_port_mapping()) - assert set(cpo.keys()) == {1} - assert cpo[1] is objs[1] - - def test_get_pluggable_obj_dict_excludes_cpo(self): - objs = self._make_obj_dict() - with patch.object(common, 'is_cpo_port', side_effect=lambda p: p in (1,)), \ - patch.object(common, 'get_port_device', side_effect=lambda p: objs[p]): - pluggable = common.get_pluggable_obj_dict(self._make_port_mapping()) - assert set(pluggable.keys()) == {0, 2} - assert pluggable[0] is objs[0] - def test_accessors_are_complementary(self): objs = self._make_obj_dict() port_mapping = self._make_port_mapping() @@ -87,6 +66,8 @@ def test_accessors_are_complementary(self): pluggable = common.get_pluggable_obj_dict(port_mapping) assert set(cpo) | set(pluggable) == set(objs) assert set(cpo) & set(pluggable) == set() + assert set(cpo) == {1} + assert cpo[1] is objs[1] def test_all_pluggable_when_no_cpo(self): objs = self._make_obj_dict() diff --git a/sonic-xcvrd/xcvrd/xcvrd_utilities/common.py b/sonic-xcvrd/xcvrd/xcvrd_utilities/common.py index 83c44425e..364e7044b 100644 --- a/sonic-xcvrd/xcvrd/xcvrd_utilities/common.py +++ b/sonic-xcvrd/xcvrd/xcvrd_utilities/common.py @@ -10,9 +10,11 @@ import subprocess import traceback import threading + from typing import Callable, Dict from swsscommon import swsscommon from sonic_py_common import syslogger, daemon_base, multi_asic from . import sfp_status_helper + from .port_event_helper import PortMapping from sonic_platform_base.sonic_xcvr.api.public.c_cmis import CmisApi except ImportError as e: @@ -137,7 +139,8 @@ def is_cpo_port(physical_port): def is_pluggable_port(physical_port): return not is_cpo_port(physical_port) -def _get_port_obj_dict(port_mapping_data, port_filter): +def _get_port_obj_dict(port_mapping_data: PortMapping | None, + port_filter: Callable[[int], bool]) -> Dict[int, object]: """ Create a dictionary mapping physical ports to their corresponding device objects, restricted to the ports accepted by port_filter. From 91a8c52e3508aa6915d283a150da8c8487647076 Mon Sep 17 00:00:00 2001 From: Brian Gallagher Date: Thu, 30 Jul 2026 00:53:33 +0000 Subject: [PATCH 08/12] Fix regressions in CMIS state machine logic Signed-off-by: Brian Gallagher --- sonic-xcvrd/tests/test_xcvrd.py | 43 +++++++++++++++++++++ sonic-xcvrd/xcvrd/cmis/cmis_manager_task.py | 15 ++++--- 2 files changed, 53 insertions(+), 5 deletions(-) diff --git a/sonic-xcvrd/tests/test_xcvrd.py b/sonic-xcvrd/tests/test_xcvrd.py index d3d7254f8..ac186e63d 100644 --- a/sonic-xcvrd/tests/test_xcvrd.py +++ b/sonic-xcvrd/tests/test_xcvrd.py @@ -2804,6 +2804,49 @@ def test_DaemonXcvrd_run_creates_cpo_and_pluggable_tasks(self, mock_init): 'CpoManagerTask', 'CpoDomInfoUpdateTask', 'CpoStateUpdateTask'): assert expected in names + def test_CmisManagerTask_on_port_update_event_without_index(self): + port_mapping = PortMapping() + for lport, pport in (('Ethernet0', 1), ('Ethernet8', 2)): + port_mapping.handle_port_change_event( + PortChangeEvent(lport, pport, 0, PortChangeEvent.PORT_ADD)) + + stop_event = threading.Event() + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) + task.force_cmis_reinit = MagicMock() + + task.on_port_update_event(PortChangeEvent('Ethernet0', -1, 0, PortChangeEvent.PORT_SET, + {'host_tx_ready': 'true'})) + assert task.port_dict['Ethernet0']['host_tx_ready'] == 'true' + assert task.force_cmis_reinit.call_count == 1 + + task.on_port_update_event(PortChangeEvent('Ethernet8', -1, 0, PortChangeEvent.PORT_SET, + {'host_tx_ready': 'true'})) + assert 'host_tx_ready' not in task.port_dict.get('Ethernet8', {}) + assert task.force_cmis_reinit.call_count == 1 + + task.on_port_update_event(PortChangeEvent('Ethernet996', -1, 0, PortChangeEvent.PORT_SET, + {'host_tx_ready': 'true'})) + assert 'Ethernet996' not in task.port_dict + assert task.force_cmis_reinit.call_count == 1 + + def test_CmisManagerTask_on_port_update_event_after_breakout(self): + port_mapping = PortMapping() + port_mapping.handle_port_change_event( + PortChangeEvent('Ethernet0', 1, 0, PortChangeEvent.PORT_ADD)) + + stop_event = threading.Event() + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) + task.force_cmis_reinit = MagicMock() + + task.on_port_update_event(PortChangeEvent('Ethernet2', 1, 0, PortChangeEvent.PORT_SET, + {'speed': '200000', 'lanes': '3,4', + 'subport': '2'})) + assert task.port_dict['Ethernet2']['index'] == 1 + + task.on_port_update_event(PortChangeEvent('Ethernet2', -1, 0, PortChangeEvent.PORT_SET, + {'host_tx_ready': 'true'})) + assert task.port_dict['Ethernet2']['host_tx_ready'] == 'true' + def test_SffManagerTask_handle_port_change_event(self): stop_event = threading.Event() task = SffManagerTask(DEFAULT_NAMESPACE, stop_event, MagicMock(), helper_logger) diff --git a/sonic-xcvrd/xcvrd/cmis/cmis_manager_task.py b/sonic-xcvrd/xcvrd/cmis/cmis_manager_task.py index b5e191469..dc235ef57 100644 --- a/sonic-xcvrd/xcvrd/cmis/cmis_manager_task.py +++ b/sonic-xcvrd/xcvrd/cmis/cmis_manager_task.py @@ -52,8 +52,13 @@ def __init__(self, namespaces, port_mapping, port_obj_dict, main_thread_stop_eve self.exc = None self.task_stopping_event = threading.Event() self.main_thread_stop_event = main_thread_stop_event - self.port_dict = {k: {"asic_id": v} for k, v in port_mapping.logical_to_asic.items()} - self.port_mapping = port_mapping + self.port_dict = {} + for lport, asic_id in port_mapping.logical_to_asic.items(): + entry = {"asic_id": asic_id} + pports = port_mapping.get_logical_to_physical(lport) + if pports: + entry["index"] = pports[0] + self.port_dict[lport] = entry self.decomm_pending_dict = {} self.isPortInitDone = False self.isPortConfigDone = False @@ -119,7 +124,8 @@ def on_port_update_event(self, port_change_event): if port_change_event.port_dict is None: return - if pport not in self.port_obj_dict: + owner_pport = pport if pport >= 0 else self.port_dict.get(lport, {}).get('index') + if owner_pport not in self.port_obj_dict: return if port_change_event.event_type == port_change_event.PORT_SET: @@ -1373,8 +1379,7 @@ def run(self): self.wait_for_port_config_done(namespace) for lport in self.port_dict.keys(): - pports = self.port_mapping.get_logical_to_physical(lport) - if pports and pports[0] in self.port_obj_dict: + if self.port_dict[lport].get('index') in self.port_obj_dict: self.update_port_transceiver_status_table_sw_cmis_state(lport, CMIS_STATE_UNKNOWN) self.task_worker() From e8509e931204fc648f634e58c64555b35661385c Mon Sep 17 00:00:00 2001 From: Brian Gallagher Date: Fri, 31 Jul 2026 23:05:46 +0000 Subject: [PATCH 09/12] Move CPO tasks into cpo directory Signed-off-by: Brian Gallagher --- sonic-xcvrd/xcvrd/cpo/cpo_state_task.py | 12 ++++++++++++ sonic-xcvrd/xcvrd/cpo/dom.py | 10 ++++++++++ sonic-xcvrd/xcvrd/dom/dom_mgr.py | 4 ---- sonic-xcvrd/xcvrd/xcvrd.py | 10 +++------- 4 files changed, 25 insertions(+), 11 deletions(-) create mode 100644 sonic-xcvrd/xcvrd/cpo/cpo_state_task.py create mode 100644 sonic-xcvrd/xcvrd/cpo/dom.py diff --git a/sonic-xcvrd/xcvrd/cpo/cpo_state_task.py b/sonic-xcvrd/xcvrd/cpo/cpo_state_task.py new file mode 100644 index 000000000..57cbd7c9b --- /dev/null +++ b/sonic-xcvrd/xcvrd/cpo/cpo_state_task.py @@ -0,0 +1,12 @@ +#!/usr/bin/env python3 + +try: + from ..xcvrd import SfpStateUpdateTask +except ImportError as e: + raise ImportError(str(e) + " - required module not found") + + +class CpoStateUpdateTask(SfpStateUpdateTask): + def __init__(self, namespaces, port_mapping, port_obj_dict, main_thread_stop_event, sfp_error_event): + super().__init__(namespaces, port_mapping, port_obj_dict, main_thread_stop_event, sfp_error_event) + self.name = "CpoStateUpdateTask" diff --git a/sonic-xcvrd/xcvrd/cpo/dom.py b/sonic-xcvrd/xcvrd/cpo/dom.py new file mode 100644 index 000000000..6d03b75eb --- /dev/null +++ b/sonic-xcvrd/xcvrd/cpo/dom.py @@ -0,0 +1,10 @@ +#!/usr/bin/env python3 + +try: + from ..dom.dom_mgr import DomInfoUpdateTask +except ImportError as e: + raise ImportError(str(e) + " - required module not found") + + +class CpoDomInfoUpdateTask(DomInfoUpdateTask): + name = "CpoDomInfoUpdateTask" diff --git a/sonic-xcvrd/xcvrd/dom/dom_mgr.py b/sonic-xcvrd/xcvrd/dom/dom_mgr.py index 063a348cc..82c1becf2 100644 --- a/sonic-xcvrd/xcvrd/dom/dom_mgr.py +++ b/sonic-xcvrd/xcvrd/dom/dom_mgr.py @@ -530,10 +530,6 @@ def on_remove_logical_port(self, port_change_event): ]) -class CpoDomInfoUpdateTask(DomInfoUpdateTask): - name = "CpoDomInfoUpdateTask" - - class DomThermalInfoUpdateTask(DomInfoUpdateBase): name = 'DomThermalInfoUpdateTask' diff --git a/sonic-xcvrd/xcvrd/xcvrd.py b/sonic-xcvrd/xcvrd/xcvrd.py index 5b115145c..35688953c 100644 --- a/sonic-xcvrd/xcvrd/xcvrd.py +++ b/sonic-xcvrd/xcvrd/xcvrd.py @@ -28,9 +28,10 @@ from .xcvrd_utilities import sfp_status_helper from .sff_mgr import SffManagerTask - from .dom.dom_mgr import DomThermalInfoUpdateTask, DomInfoUpdateTask, CpoDomInfoUpdateTask + from .dom.dom_mgr import DomThermalInfoUpdateTask, DomInfoUpdateTask from .cmis.cmis_manager_task import CmisManagerTask from .cpo.cpo_manager_task import CpoManagerTask + from .cpo.dom import CpoDomInfoUpdateTask from .xcvrd_utilities.xcvr_table_helper import * from .xcvrd_utilities import port_event_helper from .xcvrd_utilities.port_event_helper import PortChangeObserver @@ -870,12 +871,6 @@ def update_log_level(self): return self.logger.update_log_level() -class CpoStateUpdateTask(SfpStateUpdateTask): - def __init__(self, namespaces, port_mapping, port_obj_dict, main_thread_stop_event, sfp_error_event): - super().__init__(namespaces, port_mapping, port_obj_dict, main_thread_stop_event, sfp_error_event) - self.name = "CpoStateUpdateTask" - - # # Daemon ======================================================================= # @@ -1182,6 +1177,7 @@ def run(self): # Start the CPO state info update thread cpo_state_update = None if self.cpo_obj_dict: + from .cpo.cpo_state_task import CpoStateUpdateTask cpo_state_update = CpoStateUpdateTask(self.namespaces, port_mapping_data, self.cpo_obj_dict, self.stop_event, self.sfp_error_event) cpo_state_update.start() self.threads.append(cpo_state_update) From 8bc4f1777974a9f28b8988bf648c2e8dcb59bc8a Mon Sep 17 00:00:00 2001 From: Brian Gallagher Date: Tue, 4 Aug 2026 11:50:36 +0000 Subject: [PATCH 10/12] Rename cpo/dom.py to cpo/dom_mgr.py Signed-off-by: Brian Gallagher --- sonic-xcvrd/xcvrd/cpo/{dom.py => dom_mgr.py} | 0 sonic-xcvrd/xcvrd/xcvrd.py | 2 +- 2 files changed, 1 insertion(+), 1 deletion(-) rename sonic-xcvrd/xcvrd/cpo/{dom.py => dom_mgr.py} (100%) diff --git a/sonic-xcvrd/xcvrd/cpo/dom.py b/sonic-xcvrd/xcvrd/cpo/dom_mgr.py similarity index 100% rename from sonic-xcvrd/xcvrd/cpo/dom.py rename to sonic-xcvrd/xcvrd/cpo/dom_mgr.py diff --git a/sonic-xcvrd/xcvrd/xcvrd.py b/sonic-xcvrd/xcvrd/xcvrd.py index 35688953c..0fdcafcf2 100644 --- a/sonic-xcvrd/xcvrd/xcvrd.py +++ b/sonic-xcvrd/xcvrd/xcvrd.py @@ -31,7 +31,7 @@ from .dom.dom_mgr import DomThermalInfoUpdateTask, DomInfoUpdateTask from .cmis.cmis_manager_task import CmisManagerTask from .cpo.cpo_manager_task import CpoManagerTask - from .cpo.dom import CpoDomInfoUpdateTask + from .cpo.dom_mgr import CpoDomInfoUpdateTask from .xcvrd_utilities.xcvr_table_helper import * from .xcvrd_utilities import port_event_helper from .xcvrd_utilities.port_event_helper import PortChangeObserver From f413b2ddf2c6c51a2e1e7d9ade7cc5d490b6d862 Mon Sep 17 00:00:00 2001 From: Brian Gallagher Date: Tue, 4 Aug 2026 18:20:53 +0000 Subject: [PATCH 11/12] update cpo tests Signed-off-by: Brian Gallagher --- sonic-xcvrd/tests/test_xcvrd.py | 48 ++++++++++++++++----------------- 1 file changed, 24 insertions(+), 24 deletions(-) diff --git a/sonic-xcvrd/tests/test_xcvrd.py b/sonic-xcvrd/tests/test_xcvrd.py index ac186e63d..d7e852eef 100644 --- a/sonic-xcvrd/tests/test_xcvrd.py +++ b/sonic-xcvrd/tests/test_xcvrd.py @@ -3221,7 +3221,7 @@ def test_SffManagerTask_task_worker(self, mock_chassis, mock_logger): def test_CmisManagerTask_update_port_transceiver_status_table_sw_cmis_state(self): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) port_change_event = PortChangeEvent('Ethernet0', 1, 0, PortChangeEvent.PORT_SET) task.on_port_update_event(port_change_event) @@ -3279,7 +3279,7 @@ def test_CmisManagerTask_handle_port_change_event(self): def test_CmisManagerTask_get_configured_freq(self, mock_table_helper): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) cfg_port_tbl = MagicMock() cfg_port_tbl.hget = MagicMock(return_value=(True, 193100)) mock_table_helper.get_cfg_port_tbl = MagicMock(return_value=cfg_port_tbl) @@ -3291,7 +3291,7 @@ def test_CmisManagerTask_get_configured_freq(self, mock_table_helper): def test_CmisManagerTask_get_configured_tx_power_from_db(self, mock_table_helper): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) cfg_port_tbl = MagicMock() cfg_port_tbl.hget = MagicMock(return_value=(True, -10)) mock_table_helper.get_cfg_port_tbl = MagicMock(return_value=cfg_port_tbl) @@ -3496,7 +3496,7 @@ def test_CmisManagerTask_task_run_stop(self, mock_chassis): port_mapping = PortMapping() stop_event = threading.Event() - cmis_manager = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + cmis_manager = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) cmis_manager.wait_for_port_config_done = MagicMock() cmis_manager.start() cmis_manager.join() @@ -3523,7 +3523,7 @@ def test_CmisManagerTask_is_decommission_required(self, active_map, desired_map, mock_xcvr_api.get_application = MagicMock(side_effect=lambda lane: active_map[lane]) port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} task.get_desired_app_map = MagicMock(return_value=desired_map) assert task.is_decommission_required(mock_xcvr_api, 'Ethernet0') == expected @@ -3543,7 +3543,7 @@ def test_CmisManagerTask_is_decommission_required_uses_active_appsel_not_staged( port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} task.get_desired_app_map = MagicMock(return_value=desired_map) @@ -3566,7 +3566,7 @@ def test_CmisManagerTask_is_decommission_required_invalid_active_appsel(self): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} task.get_desired_app_map = MagicMock(return_value=desired_map) @@ -3589,7 +3589,7 @@ def test_CmisManagerTask_is_decommission_required_missing_active_appsel_lane(sel port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} task.get_desired_app_map = MagicMock(return_value=desired_map) @@ -3600,7 +3600,7 @@ def test_CmisManagerTask_get_desired_app_map(self): mock_xcvr_api = MagicMock() port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} # Mock cfg_port_tbl @@ -3643,7 +3643,7 @@ def test_CmisManagerTask_get_desired_app_map_mixed_mode(self): mock_xcvr_api = MagicMock() port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} cfg_port_tbl = MagicMock() @@ -3683,7 +3683,7 @@ def test_CmisManagerTask_get_desired_app_map_no_matching_app(self): mock_xcvr_api = MagicMock() port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} cfg_port_tbl = MagicMock() @@ -3702,7 +3702,7 @@ def test_CmisManagerTask_get_desired_app_map_uses_gearbox_lane_count(self): mock_xcvr_api = MagicMock() port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} cfg_port_tbl = MagicMock() @@ -3732,7 +3732,7 @@ def test_CmisManagerTask_get_sibling_port_configs(self): numeric fields, or get() returning not-found.""" port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} cfg_port_tbl = MagicMock() @@ -3777,7 +3777,7 @@ def test_CmisManagerTask_get_sibling_port_configs_no_cfg_port_tbl(self): """get_sibling_port_configs returns [] when cfg_port_tbl is None.""" port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} task.xcvr_table_helper.get_cfg_port_tbl = MagicMock(return_value=None) task._gearbox_lanes_dict = {} @@ -3788,7 +3788,7 @@ def test_CmisManagerTask_get_sibling_port_configs_uses_gearbox_lane_count(self): """get_sibling_port_configs uses gearbox lane count when present.""" port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) task.port_dict['Ethernet0'] = {'index': 1, 'asic_id': 0} cfg_port_tbl = MagicMock() @@ -3859,7 +3859,7 @@ def get_application(lane): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) assert task.is_cmis_application_update_required(mock_xcvr_api, app_new, host_lanes_mask) == expected @@ -3928,7 +3928,7 @@ def get_host_lane_assignment_option_side_effect(app): mock_xcvr_api.get_host_lane_assignment_option = MagicMock(side_effect=get_host_lane_assignment_option_side_effect) port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) appl = common.get_cmis_application_desired(mock_xcvr_api, host_lane_count, speed) assert task.get_cmis_host_lanes_mask(mock_xcvr_api, appl, host_lane_count, subport) == expected @@ -4034,7 +4034,7 @@ def mock_get_side_effect(key): def test_CmisManagerTask_get_host_lane_count(self, gearbox_lanes_dict, lport, port_config_lanes, expected_count): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) task._gearbox_lanes_dict = gearbox_lanes_dict result = task.get_host_lane_count(lport, port_config_lanes) @@ -4044,7 +4044,7 @@ def test_CmisManagerTask_gearbox_integration_end_to_end(self): """Test end-to-end integration of gearbox line lanes with CMIS application selection""" port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) # Mock gearbox lanes dictionary - port has 4 system lanes but only 2 line lanes gearbox_lanes_dict = {"Ethernet0": 2} # 2 line lanes from gearbox @@ -4085,7 +4085,7 @@ def test_CmisManagerTask_gearbox_caching_integration(self): """Test that gearbox lanes dictionary is properly cached and used in task worker""" port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) # Mock the XcvrTableHelper to return a gearbox lanes dictionary mock_gearbox_lanes_dict = {"Ethernet0": 2, "Ethernet4": 4} @@ -4111,7 +4111,7 @@ def test_CmisManagerTask_post_port_active_apsel_to_db_error_cases(self, mock_fie port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) lport = "Ethernet0" host_lanes_mask = 0xff @@ -4173,7 +4173,7 @@ def test_CmisManagerTask_post_port_active_apsel_to_db(self): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) task.xcvr_table_helper = XcvrTableHelper(DEFAULT_NAMESPACE) task.xcvr_table_helper.get_intf_tbl = MagicMock(return_value=int_tbl) @@ -4270,7 +4270,7 @@ def test_CmisManagerTask_post_port_active_apsel_to_db(self): def test_CmisManagerTask_test_is_timer_expired(self, expired_time, current_time, expected_result): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) # Call the is_timer_expired function result = task.is_timer_expired(expired_time, current_time) @@ -6503,7 +6503,7 @@ def test_CmisManagerTask_validate_frequency_and_grid(self, lport, freq, grid, ex mock_xcvr_api.get_supported_freq_config.return_value = (0x80, 0, 0, 191300, 196100) port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, MagicMock(), stop_event) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) result = task.validate_frequency_and_grid(mock_xcvr_api, lport, freq, grid) assert result == expected From 2b656865edca9c58db9813d22b968e97d854ada5 Mon Sep 17 00:00:00 2001 From: Brian Gallagher Date: Wed, 5 Aug 2026 09:23:53 +0000 Subject: [PATCH 12/12] Fix new unit-tests Signed-off-by: Brian Gallagher --- sonic-xcvrd/tests/test_xcvrd.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/sonic-xcvrd/tests/test_xcvrd.py b/sonic-xcvrd/tests/test_xcvrd.py index d7c9214f6..5dea90a25 100644 --- a/sonic-xcvrd/tests/test_xcvrd.py +++ b/sonic-xcvrd/tests/test_xcvrd.py @@ -3241,7 +3241,7 @@ def test_CmisManagerTask_is_fast_reboot_enabled_for_lport(self): port_mapping = PortMapping() port_mapping.logical_to_asic = {'Ethernet0': 1} stop_event = threading.Event() - task = CmisManagerTask(['asic1'], port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(['asic1'], port_mapping, {1: MagicMock()}, stop_event) common.is_fast_reboot_enabled.reset_mock() assert task.is_fast_reboot_enabled_for_lport('Ethernet0') is True @@ -3252,7 +3252,7 @@ def test_CmisManagerTask_is_fast_reboot_enabled_for_lport(self): def test_CmisManagerTask_is_fast_reboot_enabled_for_lport_default_namespace(self): port_mapping = PortMapping() stop_event = threading.Event() - task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, stop_event, platform_chassis=MagicMock()) + task = CmisManagerTask(DEFAULT_NAMESPACE, port_mapping, {1: MagicMock()}, stop_event) # Ethernet999 not in port_dict -> asic_id -1 -> namespace '' with patch.object(common, 'get_namespace_from_asic_id', MagicMock()) as mock_get_ns: assert task.is_fast_reboot_enabled_for_lport('Ethernet999') is False