21 Commits
Author SHA1 Message Date
muellerr c9ae13c16f some fixes for common tmtc code 2023-01-13 11:07:08 +01:00
muellerr 8ad775e5be minor improvements 2022-10-19 11:55:58 +02:00
muellerr 9abfba1e0e basic printout for received PDU TMs 2022-10-19 11:49:44 +02:00
muellerr 34203461fa some changes for CFDP done logic 2022-10-17 12:26:00 +02:00
muellerr 73ba9f5d90 use singular enum names 2022-09-15 18:46:28 +02:00
muellerr 7b0c8fa25a separate function for cfdp setup 2022-09-15 16:28:58 +02:00
muellerr 4a9ed5aad0 more generic PDU packing approach 2022-09-15 14:31:13 +02:00
muellerr 51871a887f moved a lot of code to tmtccmd 2022-09-14 19:01:47 +02:00
muellerr 3787164cce this works 2022-09-14 18:30:52 +02:00
muellerr c4dcf0df76 getting complicated 2022-09-14 18:19:56 +02:00
muellerr 061d62579b info printout for finished CFDP queue 2022-09-14 17:54:57 +02:00
muellerr 8237591c4c more useful printout 2022-09-14 17:48:15 +02:00
muellerr e4e9b3d5d3 all indications working 2022-09-14 17:40:39 +02:00
muellerr 0faccfb456 better printout 2022-09-14 17:35:33 +02:00
muellerr 0e1d451651 send eof packet to finish transaction 2022-09-14 14:00:37 +02:00
muellerr 6ee3fa67dc send wrapped CFDP packets and clean up printout 2022-09-13 13:54:47 +02:00
muellerr a3d1808457 building infastructure to build queue 2022-09-13 11:59:29 +02:00
muellerr e9aaa5502d printout tweaks 2022-09-12 10:08:03 +02:00
muellerr 845ee041d4 correct code is now reached 2022-09-12 09:39:37 +02:00
muellerr 99ccd04b2a some adaptions 2022-09-12 00:53:15 +02:00
muellerr 27c1002d25 now only arg param to proc conversion missing 2022-09-09 17:51:41 +02:00
3 changed files with 123 additions and 198 deletions
+117 -193
View File
@@ -1,32 +1,30 @@
import logging
import sys
from pathlib import Path
from typing import Optional, Sequence
from typing import cast
from spacepackets import SpacePacket, SpacePacketHeader, PacketTypes
from spacepackets import SpacePacket, SpacePacketHeader
from spacepackets.ccsds import SPACE_PACKET_HEADER_SIZE
from spacepackets.cfdp import (
TransmissionModes,
PduType,
DirectiveType,
GenericPduPacket,
PduHolder,
PduFactory,
ChecksumTypes,
TransmissionMode,
ChecksumType,
ConditionCode,
PduHolder,
DirectiveType,
PduFactory,
PduType,
)
from spacepackets.cfdp.pdu import MetadataPdu, FileDataPdu
from tmtccmd.cfdp import (
RemoteEntityCfgTable,
RemoteEntityCfg,
LocalEntityCfg,
CfdpUserBase,
IndicationCfg,
TransactionId,
)
from spacepackets.util import UnsignedByteField
from tmtccmd.cfdp.defs import CfdpRequestType, CfdpStates
from tmtccmd.cfdp.handler import SourceHandler, DestHandler
from tmtccmd.cfdp.defs import CfdpRequestType
from tmtccmd.cfdp.handler import CfdpInCcsdsHandler
from tmtccmd.cfdp.mib import DefaultFaultHandlerBase
from tmtccmd.cfdp.request import PutRequest, PutRequestCfg
from tmtccmd.cfdp.user import (
FileSegmentRecvdParams,
MetadataRecvParams,
@@ -36,7 +34,7 @@ from tmtccmd.config.args import ProcedureParamsWrapper
from tmtccmd.logging import get_current_time_string
from tmtccmd.pus.pus_11_tc_sched import Subservices as Pus11Subservices
from tmtccmd.tc.queue import DefaultPusQueueHelper
from tmtccmd.util import FileSeqCountProvider, PusFileSeqCountProvider, ProvidesSeqCount
from tmtccmd.util import FileSeqCountProvider, PusFileSeqCountProvider
from tmtccmd.util.tmtc_printer import FsfwTmTcPrinter
try:
@@ -72,12 +70,11 @@ from tmtccmd import (
TcHandlerBase,
get_console_logger,
TmTcCfgHookBase,
BackendBase,
DefProcedureParams,
CcsdsTmtcBackend,
)
from tmtccmd.pus import VerificationWrapper
from tmtccmd.tc import (
ProcedureHelper,
ProcedureWrapper,
FeedWrapper,
TcProcedureType,
TcQueueEntryType,
@@ -91,6 +88,7 @@ from tmtccmd.config import (
SetupWrapper,
SetupParams,
PreArgsParsingWrapper,
params_to_procedure_conversion,
)
from common_tmtc.pus_tm.factory_hook import pus_factory_hook
@@ -114,13 +112,13 @@ class ExampleCfdpFaultHandler(DefaultFaultHandlerBase):
class ExampleCfdpUser(CfdpUserBase):
def transaction_indication(self, transaction_id: TransactionId):
pass
LOGGER.info(f"CFDP User: Start of File {transaction_id}")
def eof_sent_indication(self, transaction_id: TransactionId):
pass
LOGGER.info(f"CFDP User: EOF sent for {transaction_id}")
def transaction_finished_indication(self, params: TransactionFinishedParams):
pass
LOGGER.info(f"CFDP User: {params.transaction_id} finished")
def metadata_recv_indication(self, params: MetadataRecvParams):
pass
@@ -153,137 +151,26 @@ class ExampleCfdpUser(CfdpUserBase):
pass
class CfdpCcsdsWrapper(SpecificApidHandlerBase):
def __init__(
self,
cfg: LocalEntityCfg,
user: CfdpUserBase,
remote_cfgs: Sequence[RemoteEntityCfg],
ccsds_apid: int,
cfdp_seq_cnt_provider: ProvidesSeqCount,
ccsds_seq_cnt_provider: ProvidesSeqCount,
):
"""Wrapper helper type used to wrap PDU packets into CCSDS packets and to extract PDU
packets from CCSDS packets.
:param cfg: Local CFDP entity configuration.
:param user: User wrapper. This contains the indication callback implementations and the
virtual filestore implementation.
:param cfdp_seq_cnt_provider: Every CFDP file transfer has a transaction sequence number.
This provider is used to retrieve that sequence number.
:param ccsds_seq_cnt_provider: Each CFDP PDU is wrapped into a CCSDS space packet, and each
space packet has a dedicated sequence count. This provider is used to retrieve the
sequence count.
:param ccsds_apid: APID to use for the CCSDS space packet header wrapped around each PDU.
This is important so that the OBSW can distinguish between regular PUS packets and
CFDP packets.
"""
class CfdpInCcsdsWrapper(SpecificApidHandlerBase):
def __init__(self, cfdp_in_ccsds_handler: CfdpInCcsdsHandler):
super().__init__(EXAMPLE_CFDP_APID, None)
self.handler = CfdpHandler(cfg, user, cfdp_seq_cnt_provider, remote_cfgs)
self.ccsds_seq_cnt_provider = ccsds_seq_cnt_provider
self.ccsds_apid = ccsds_apid
def pull_next_dest_packet(self) -> Optional[SpacePacket]:
"""Retrieves the next PDU to send and wraps it into a space packet"""
next_packet = self.handler.pull_next_dest_packet()
if next_packet is None:
return next_packet
sp_header = SpacePacketHeader(
packet_type=PacketTypes.TC,
apid=self.ccsds_apid,
seq_count=self.ccsds_seq_cnt_provider.get_and_increment(),
data_len=next_packet.packet_len - 1,
)
return SpacePacket(sp_header, None, next_packet.pack())
def confirm_dest_packet_sent(self):
self.handler.confirm_dest_packet_sent()
def pass_packet(self, space_packet: SpacePacket):
# Unwrap the user data and pass it to the handler
pdu_raw = space_packet.user_data
pdu_base = PduFactory.from_raw(pdu_raw)
self.handler.pass_packet(pdu_base)
self.handler = cfdp_in_ccsds_handler
def handle_tm(self, packet: bytes, _user_args: any):
ccsds_header_raw = packet[0:6]
sp_header = SpacePacketHeader.unpack(ccsds_header_raw)
pdu = packet[6:]
sp = SpacePacket(sp_header, sec_header=None, user_data=pdu)
self.pass_packet(sp)
class CfdpHandler:
def __init__(
self,
cfg: LocalEntityCfg,
user: CfdpUserBase,
seq_cnt_provider: ProvidesSeqCount,
remote_cfgs: Sequence[RemoteEntityCfg],
):
self.dest_id = UnsignedByteField(EXAMPLE_PUS_APID, 2)
self.remote_cfg_table = RemoteEntityCfgTable()
self.remote_cfg_table.add_configs(remote_cfgs)
self.dest_handler = DestHandler(cfg, user, self.remote_cfg_table)
self.source_handler = SourceHandler(cfg, seq_cnt_provider, user)
def put_request(self, request: PutRequest):
if self.remote_cfg_table.get_cfg(request.cfg.destination_id):
raise ValueError(
f"No remote CFDP config found for entity ID {request.cfg.destination_id}"
)
self.source_handler.put_request(
request, self.remote_cfg_table.get_cfg(self.dest_id)
)
def pull_next_dest_packet(self) -> Optional[PduHolder]:
res = self.dest_handler.state_machine()
if res.states.packet_ready:
return self.dest_handler.pdu_holder
return None
def put_request_pending(self) -> bool:
return self.source_handler.states.state != CfdpStates.IDLE
def confirm_dest_packet_sent(self):
self.dest_handler.confirm_packet_sent_advance_fsm()
def pass_packet(self, packet: GenericPduPacket):
"""This function routes the packets based on PDU type and directive type if applicable.
The routing is based on section 4.5 of the CFDP standard whcih specifies the PDU forwarding
procedure.
"""
if packet.pdu_type == PduType.FILE_DATA:
self.dest_handler.pass_packet(packet)
# Ignore the space packet header. Its only purpose is to use the same protocol and
# have a seaprate APID for space packets. If this function is called, the APID is correct.
pdu = packet[SPACE_PACKET_HEADER_SIZE:]
pdu_base = PduFactory.from_raw(pdu)
if pdu_base.pdu_type == PduType.FILE_DATA:
LOGGER.info("Received File Data PDU TM")
else:
if packet.directive_type in [
DirectiveType.METADATA_PDU,
DirectiveType.EOF_PDU,
DirectiveType.PROMPT_PDU,
]:
# Section b) of 4.5.3: These PDUs should always be targeted towards the file
# receiver a.k.a. the destination handler
self.dest_handler.pass_packet(packet)
elif packet.directive_type in [
DirectiveType.FINISHED_PDU,
DirectiveType.NAK_PDU,
DirectiveType.KEEP_ALIVE_PDU,
]:
# Section c) of 4.5.3: These PDUs should always be targeted towards the file sender
# a.k.a. the source handler
self.source_handler.pass_packet(packet)
elif packet.directive_type == DirectiveType.ACK_PDU:
# Section a): Recipient depends on the type of PDU that is being acknowledged.
# We can simply extract the PDU type from the raw stream. If it is an EOF PDU,
# this packet is passed to the source handler. For a finished PDU, it is
# passed to the destination handler
pdu_holder = PduHolder(packet)
ack_pdu = pdu_holder.to_ack_pdu()
if ack_pdu.directive_code_of_acked_pdu == DirectiveType.EOF_PDU:
self.source_handler.pass_packet(packet)
elif ack_pdu.directive_code_of_acked_pdu == DirectiveType.FINISHED_PDU:
self.dest_handler.pass_packet(packet)
if pdu_base.directive_type == DirectiveType.FINISHED_PDU:
LOGGER.info(f"Received Finished PDU TM")
else:
LOGGER.info(
f"Received File Directive PDU with type {pdu_base.directive_type!r} TM"
)
self.handler.pass_pdu_packet(pdu_base)
class PusHandler(SpecificApidHandlerBase):
@@ -311,12 +198,14 @@ class TcHandler(TcHandlerBase):
def __init__(
self,
seq_count_provider: FileSeqCountProvider,
cfdp_in_ccsds_wrapper: CfdpCcsdsWrapper,
cfdp_in_ccsds_wrapper: CfdpInCcsdsWrapper,
pus_verificator: PusVerificator,
file_logger: logging.Logger,
raw_logger: RawTmtcTimedLogWrapper,
):
super().__init__()
self.cfdp_handler_started = False
self.cfdp_dest_id = CFDP_REMOTE_ENTITY_ID
self.seq_count_provider = seq_count_provider
self.pus_verificator = pus_verificator
self.file_logger = file_logger
@@ -329,14 +218,17 @@ class TcHandler(TcHandlerBase):
)
self.cfdp_in_ccsds_wrapper = cfdp_in_ccsds_wrapper
def feed_cb(self, info: ProcedureHelper, wrapper: FeedWrapper):
def cfdp_done(self) -> bool:
return not self.cfdp_in_ccsds_wrapper.handler.put_request_pending()
def feed_cb(self, info: ProcedureWrapper, wrapper: FeedWrapper):
self.queue_helper.queue_wrapper = wrapper.queue_wrapper
if info.proc_type == TcProcedureType.DEFAULT:
self.handle_default_procedure(info)
elif info.proc_type == TcProcedureType.CFDP:
self.handle_cfdp_procedure(info)
def handle_default_procedure(self, info: ProcedureHelper):
def handle_default_procedure(self, info: ProcedureWrapper):
def_proc = info.to_def_procedure()
service = def_proc.service
op_code = def_proc.op_code
@@ -358,33 +250,52 @@ class TcHandler(TcHandlerBase):
return pack_service_200_commands_into(q=self.queue_helper, op_code=op_code)
LOGGER.warning("Invalid Service !")
def handle_cfdp_procedure(self, info: ProcedureHelper):
def handle_cfdp_procedure(self, info: ProcedureWrapper):
cfdp_procedure = info.to_cfdp_procedure()
if cfdp_procedure.cfdp_request_type == CfdpRequestType.PUT:
if not self.cfdp_in_ccsds_wrapper.handler.put_request_pending():
# put_req = cfdp_procedure.request_wrapper.to_put_request()
with open("/tmp/hello.txt") as of:
of.write("Hello first CFDP file")
# TODO: Replace hardcoded put request if the request works
put_request_cfg = PutRequestCfg(
destination_id=self.cfdp_in_ccsds_wrapper.handler.dest_id,
source_file=Path("/tmp/hello.txt"),
dest_file="/tmp/hello-created-by-dest.txt",
closure_requested=False,
trans_mode=TransmissionModes.UNACKNOWLEDGED,
if (
not self.cfdp_in_ccsds_wrapper.handler.put_request_pending()
and not self.cfdp_handler_started
):
put_req = cfdp_procedure.request_wrapper.to_put_request()
put_req.cfg.destination_id = self.cfdp_dest_id
LOGGER.info(
f"CFDP: Starting file put request with parameters:\n{put_req}"
)
put_req = PutRequest(put_request_cfg)
LOGGER.info(f"Starting put request {put_req}")
# TODO: Only start put request if there isn't one pending yet. The source handler
# state can probably be used for this.
self.cfdp_in_ccsds_wrapper.handler.put_request(put_req)
pass
pass
self.cfdp_in_ccsds_wrapper.handler.cfdp_handler.put_request(put_req)
self.cfdp_handler_started = True
for source_pair, dest_pair in self.cfdp_in_ccsds_wrapper.handler:
pdu, sp = source_pair
pdu = cast(PduHolder, pdu)
if pdu.is_file_directive:
if pdu.pdu_directive_type == DirectiveType.METADATA_PDU:
metadata = pdu.to_metadata_pdu()
self.queue_helper.add_log_cmd(
f"CFDP Source: Sending Metadata PDU for file with size "
f"{metadata.file_size}"
)
elif pdu.pdu_directive_type == DirectiveType.EOF_PDU:
self.queue_helper.add_log_cmd(
f"CFDP Source: Sending EOF PDU"
)
else:
fd_pdu = pdu.to_file_data_pdu()
self.queue_helper.add_log_cmd(
f"CFDP Source: Sending File Data PDU for segment at offset "
f"{fd_pdu.offset} with length {len(fd_pdu.file_data)}"
)
self.queue_helper.add_ccsds_tc(sp)
self.cfdp_in_ccsds_wrapper.handler.confirm_source_packet_sent()
self.cfdp_in_ccsds_wrapper.handler.source_handler.state_machine()
def send_cb(self, params: SendCbParams):
if params.entry.is_tc:
if params.entry.entry_type == TcQueueEntryType.PUS_TC:
self.handle_tc_send_cb(params)
elif params.entry.entry_type == TcQueueEntryType.CCSDS_TC:
cfdp_packet_in_ccsds = params.entry.to_space_packet_entry()
params.com_if.send(cfdp_packet_in_ccsds.space_packet.pack())
elif params.entry.entry_type == TcQueueEntryType.LOG:
log_entry = params.entry.to_log_entry()
LOGGER.info(log_entry.log_str)
@@ -407,17 +318,19 @@ class TcHandler(TcHandlerBase):
self.file_logger.info(f"{get_current_time_string(True)}: {tc_info_string}")
params.com_if.send(raw_tc)
def queue_finished_cb(self, info: ProcedureHelper):
if info is not None and info.proc_type == TcQueueEntryType.PUS_TC:
def_proc = info.to_def_procedure()
LOGGER.info(
f"Finished queue for service {def_proc.service} and op code {def_proc.op_code}"
)
def queue_finished_cb(self, info: ProcedureWrapper):
if info is not None:
if info.proc_type == TcQueueEntryType.PUS_TC:
def_proc = info.to_def_procedure()
LOGGER.info(
f"Finished queue for service {def_proc.service} and op code {def_proc.op_code}"
)
elif info.proc_type == TcProcedureType.CFDP:
LOGGER.info(f"Finished CFDP queue")
self.cfdp_sending_done = True
def setup_params(
hook_obj: TmTcCfgHookBase, proc_param_wrapper: ProcedureParamsWrapper
) -> SetupWrapper:
def setup_params(hook_obj: TmTcCfgHookBase) -> SetupWrapper:
print(f"-- eive TMTC Commander --")
print(f"-- spacepackets v{spacepackets.__version__} --")
params = SetupParams()
@@ -428,21 +341,19 @@ def setup_params(
post_arg_parsing_wrapper = parser_wrapper.parse(hook_obj)
tmtccmd.init_printout(post_arg_parsing_wrapper.use_gui)
use_prompts = not post_arg_parsing_wrapper.use_gui
proc_param_wrapper = ProcedureParamsWrapper()
if use_prompts:
post_arg_parsing_wrapper.set_params_with_prompts(params, proc_param_wrapper)
else:
post_arg_parsing_wrapper.set_params_without_prompts(params, proc_param_wrapper)
params.apid = EXAMPLE_PUS_APID
setup_wrapper = SetupWrapper(hook_obj=hook_obj, setup_params=params)
setup_wrapper = SetupWrapper(
hook_obj=hook_obj, setup_params=params, proc_param_wrapper=proc_param_wrapper
)
return setup_wrapper
def setup_tmtc_handlers(
verif_wrapper: VerificationWrapper,
printer: FsfwTmTcPrinter,
raw_logger: RawTmtcTimedLogWrapper,
) -> (CcsdsTmHandler, TcHandler):
def setup_cfdp_handler() -> CfdpInCcsdsWrapper:
fh_base = ExampleCfdpFaultHandler()
cfdp_cfg = LocalEntityCfg(
local_entity_id=CFDP_LOCAL_ENTITY_ID,
@@ -455,8 +366,8 @@ def setup_tmtc_handlers(
max_file_segment_len=1024,
check_limit=None,
crc_on_transmission=False,
crc_type=ChecksumTypes.CRC_32,
default_transmission_mode=TransmissionModes.UNACKNOWLEDGED,
crc_type=ChecksumType.CRC_32,
default_transmission_mode=TransmissionMode.UNACKNOWLEDGED,
)
cfdp_seq_count_provider = FileSeqCountProvider(
max_bit_width=16, file_name=Path("seqcnt_cfdp_transaction.txt")
@@ -465,7 +376,7 @@ def setup_tmtc_handlers(
file_name=Path("seqcnt_cfdp_ccsds_.txt")
)
cfdp_user = ExampleCfdpUser()
cfdp_in_ccsds_handler = CfdpCcsdsWrapper(
cfdp_in_ccsds_handler = CfdpInCcsdsHandler(
cfg=cfdp_cfg,
remote_cfgs=[remote_cfg],
ccsds_apid=EXAMPLE_CFDP_APID,
@@ -473,12 +384,21 @@ def setup_tmtc_handlers(
cfdp_seq_cnt_provider=cfdp_seq_count_provider,
user=cfdp_user,
)
return CfdpInCcsdsWrapper(cfdp_in_ccsds_handler)
def setup_tmtc_handlers(
verif_wrapper: VerificationWrapper,
printer: FsfwTmTcPrinter,
raw_logger: RawTmtcTimedLogWrapper,
) -> (CcsdsTmHandler, TcHandler):
cfdp_in_ccsds_wrapper = setup_cfdp_handler()
pus_handler = PusHandler(
printer=printer, raw_logger=raw_logger, wrapper=verif_wrapper
)
ccsds_handler = CcsdsTmHandler(None)
ccsds_handler.add_apid_handler(pus_handler)
ccsds_handler.add_apid_handler(cfdp_in_ccsds_wrapper)
tc_handler = TcHandler(
file_logger=printer.file_logger,
raw_logger=raw_logger,
@@ -486,7 +406,7 @@ def setup_tmtc_handlers(
seq_count_provider=PusFileSeqCountProvider(
file_name=Path("seqcnt_pus_ccsds.txt")
),
cfdp_in_ccsds_wrapper=cfdp_in_ccsds_handler,
cfdp_in_ccsds_wrapper=cfdp_in_ccsds_wrapper,
)
return ccsds_handler, tc_handler
@@ -495,9 +415,13 @@ def setup_backend(
setup_wrapper: SetupWrapper,
tc_handler: TcHandler,
ccsds_handler: CcsdsTmHandler,
) -> BackendBase:
) -> CcsdsTmtcBackend:
init_proc = params_to_procedure_conversion(setup_wrapper.proc_param_wrapper)
tmtc_backend = tmtccmd.create_default_tmtc_backend(
setup_wrapper=setup_wrapper, tm_handler=ccsds_handler, tc_handler=tc_handler
setup_wrapper=setup_wrapper,
tm_handler=ccsds_handler,
tc_handler=tc_handler,
init_procedure=init_proc,
)
tmtccmd.start(tmtc_backend=tmtc_backend, hook_obj=setup_wrapper.hook_obj)
return tmtc_backend
return cast(CcsdsTmtcBackend, tmtc_backend)
+5 -3
View File
@@ -1,12 +1,14 @@
from tmtccmd.pus.pus_17_test import (
pack_service_17_ping_command,
pack_generic_service17_test,
)
from tmtccmd.tc.queue import DefaultPusQueueHelper
from tmtccmd.logging import get_console_logger
LOGGER = get_console_logger()
def pack_service_17_commands(op_code: str, q: DefaultPusQueueHelper):
if op_code == "0":
if op_code in ["0", "ping"]:
q.add_pus_tc(pack_service_17_ping_command())
else:
pack_generic_service17_test(q=q)
LOGGER.warning(f"Invalid op code {op_code}")
+1 -2
View File
@@ -1,8 +1,7 @@
from datetime import timedelta
from spacepackets.ecss.tc import PusTelecommand
from deps.spacepackets.spacepackets.ecss import PusServices
from spacepackets.ecss import PusServices
from tmtccmd.config import TmtcDefinitionWrapper, OpCodeEntry
from tmtccmd.tc.pus_200_fsfw_modes import pack_mode_data, Modes
from tmtccmd.tc.pus_20_params import (