Back to home page

EIC code displayed by LXR

 
 

    


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 #: Payload keys carried into event payloads when present.
0017 PAYLOAD_KEYS = ('filename', 'sequence', 'site', 'req_id', 'input_dataset',
0018                 'dataset', 'state', 'substate', 'reason', 'container')
0019 
0020 #: Testbed agents stamp bus messages with naive local time; the agents
0021 #: run in this zone.
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]