File indexing completed on 2026-08-12 09:36:21
0001 """Shared message-to-event mapping for testbed workflow episodes.
0002
0003 Testbed agents stamp every bus message with sender, msg_type,
0004 timestamp, namespace, execution_id, and run_id; the mapping here turns
0005 that shape into episode events and participant sightings without any
0006 workflow-specific knowledge. Workflow definitions subclass
0007 ``TestbedEpisodeDefinition`` and add their completion passes.
0008 """
0009
0010 from datetime import datetime
0011 from typing import Dict, List, Optional
0012 from zoneinfo import ZoneInfo
0013
0014 from swf_common_lib.episodes import EpisodeDefinition, utc_now_iso
0015
0016
0017 PAYLOAD_KEYS = ('filename', 'sequence', 'site', 'req_id', 'input_dataset',
0018 'dataset', 'state', 'substate', 'reason', 'container')
0019
0020
0021
0022 AGENT_ZONE = ZoneInfo('America/New_York')
0023
0024
0025 def agent_kind(sender: str) -> str:
0026 """Sender agent name -> participant kind (e.g. 'daq_simulator')."""
0027 return sender.rsplit('-agent-', 1)[0] if '-agent-' in sender else 'agent'
0028
0029
0030 def message_time(message: Dict) -> Optional[str]:
0031 """A message's timestamp as timezone-aware ISO, or None.
0032
0033 Bus timestamps are naive local stamps from the sending agent; the
0034 store accepts only aware times, so the agents' zone is attached.
0035 Writers that stamp nothing (a known gap) fall back to the recorded
0036 sent time (backfill) or the builder's arrival stamp (live).
0037 """
0038 for field in ('timestamp', 'sent_at', '_received_at'):
0039 raw = message.get(field)
0040 if not raw:
0041 continue
0042 try:
0043 parsed = datetime.fromisoformat(str(raw))
0044 except ValueError:
0045 continue
0046 if parsed.tzinfo is None:
0047 parsed = parsed.replace(tzinfo=AGENT_ZONE)
0048 return parsed.isoformat()
0049 return None
0050
0051
0052 class TestbedEpisodeDefinition(EpisodeDefinition):
0053 """Message mapping shared by all testbed workflow definitions."""
0054
0055 scope = 'testbed'
0056
0057 def started_at(self, message: Dict) -> str:
0058 return message_time(message) or utc_now_iso()
0059
0060 def ended_at(self, message: Dict) -> str:
0061 return message_time(message) or utc_now_iso()
0062
0063 def event_from_message(self, message: Dict) -> Optional[Dict]:
0064 sender = message.get('sender') or message.get('sender_agent')
0065 msg_type = message.get('msg_type')
0066 when = message_time(message)
0067 if not (sender and msg_type and when):
0068 return None
0069 payload = {key: message[key] for key in PAYLOAD_KEYS
0070 if message.get(key) is not None}
0071 return {'time': when, 'kind': msg_type, 'participant': sender,
0072 'payload': payload}
0073
0074 def participants_from_message(self, message: Dict) -> List[Dict]:
0075 sender = message.get('sender') or message.get('sender_agent')
0076 when = message_time(message)
0077 if not sender:
0078 return []
0079 entry = {'id': sender, 'label': sender, 'kind': agent_kind(sender)}
0080 if when:
0081 entry['born_at'] = when
0082 return [entry]