@@ -121,15 +121,21 @@ def generate(self) -> None:
121121 self ._generate_baseline ()
122122 self ._report_progress ("phase_end" , {"phase" : "baseline" })
123123
124- # Phase 3: Execute storyline events (if present)
124+ # Phase 6. 3: Execute remaining storyline events not covered by baseline hours
125125 if self .scenario .storyline :
126- logger .info (f"Executing { len (self .scenario .storyline )} storyline events" )
127- self ._report_progress ("phase_start" , {
128- "phase" : "storyline" ,
129- "description" : f"Executing { len (self .scenario .storyline )} storyline events"
130- })
131- self ._execute_storyline ()
132- self ._report_progress ("phase_end" , {"phase" : "storyline" })
126+ remaining = [i for i in range (len (self .scenario .storyline ))
127+ if i not in self ._storyline_executed ]
128+ if remaining :
129+ logger .info (f"Executing { len (remaining )} remaining storyline events (outside baseline window)" )
130+ self ._report_progress ("phase_start" , {
131+ "phase" : "storyline" ,
132+ "description" : f"Executing { len (remaining )} remaining storyline events"
133+ })
134+ for idx in remaining :
135+ self ._execute_single_storyline_event (idx )
136+ self ._storyline_executed .add (idx )
137+ self ._barrier_flush_all_emitters ()
138+ self ._report_progress ("phase_end" , {"phase" : "storyline" })
133139 finally :
134140 # Phase 4: Finalize and close emitters (always, even on error)
135141 self ._report_progress ("phase_start" , {"phase" : "finalize" , "description" : "Finalizing generation" })
@@ -255,6 +261,20 @@ def _initialize(self) -> None:
255261 self ._system_pids : dict [str , dict [str , int ]] = {} # hostname -> {role: pid}
256262 self ._seed_system_process_trees ()
257263
264+ # Phase 6.3: Pre-parse storyline event times for interleaved generation
265+ self ._storyline_by_hour : dict [int , list ] = {} # hour_epoch -> list of (time, event_idx)
266+ if self .scenario .storyline :
267+ for idx , event in enumerate (self .scenario .storyline ):
268+ event_time = self ._parse_storyline_time (event .time )
269+ hour_key = int (event_time .replace (minute = 0 , second = 0 , microsecond = 0 ).timestamp ())
270+ self ._storyline_by_hour .setdefault (hour_key , []).append ((event_time , idx ))
271+ # Sort each hour's events by time
272+ for key in self ._storyline_by_hour :
273+ self ._storyline_by_hour [key ].sort ()
274+ logger .info (f"Pre-parsed { len (self .scenario .storyline )} storyline events across { len (self ._storyline_by_hour )} hours" )
275+
276+ self ._storyline_executed : set [int ] = set ()
277+
258278 logger .info ("Initialization complete" )
259279
260280 def _generate_baseline (self ) -> None :
@@ -326,6 +346,13 @@ def _generate_baseline(self) -> None:
326346 # Phase 5.4: Generate system traffic (DNS, NTP, scheduled tasks)
327347 self ._generate_system_traffic (current_hour )
328348
349+ # Phase 6.3: Interleave storyline events into this hour
350+ hour_key = int (current_hour .timestamp ())
351+ for event_time , event_idx in self ._storyline_by_hour .get (hour_key , []):
352+ if event_idx not in self ._storyline_executed :
353+ self ._execute_single_storyline_event (event_idx )
354+ self ._storyline_executed .add (event_idx )
355+
329356 # Phase 5.2: Terminate stale processes
330357 self ._terminate_stale_processes (current_hour )
331358
@@ -854,6 +881,39 @@ def _execute_storyline(self) -> None:
854881 # Barrier flush after each storyline event (ensures event written before proceeding)
855882 self ._barrier_flush_all_emitters ()
856883
884+ def _execute_single_storyline_event (self , event_idx : int ) -> None :
885+ """Execute a single storyline event by index (used for interleaved generation)."""
886+ storyline_event = self .scenario .storyline [event_idx ]
887+ event_num = event_idx + 1
888+
889+ # Parse event time with jitter
890+ event_time = self ._parse_storyline_time (storyline_event .time )
891+ jitter_rng = random .Random (hash (f"jitter_{ event_num } _{ self .scenario .name } " ))
892+ jitter = timedelta (
893+ seconds = jitter_rng .uniform (- 30 , 30 ),
894+ microseconds = jitter_rng .randint (0 , 999999 ),
895+ )
896+ event_time = event_time + jitter
897+
898+ actor = self ._find_actor (storyline_event .actor )
899+ system = self ._find_system (storyline_event .system )
900+ if not actor or not system :
901+ return
902+
903+ logger .info (f"Executing interleaved storyline event: { storyline_event .actor } on { storyline_event .system } at { event_time } " )
904+
905+ event_types = self ._match_activity_to_events (storyline_event .activity )
906+ self .state_manager .set_current_time (event_time )
907+
908+ for event_type in event_types :
909+ malicious_event = self ._execute_storyline_event (
910+ actor = actor , system = system , time = event_time ,
911+ event_type = event_type , activity = storyline_event .activity ,
912+ details = storyline_event .details ,
913+ )
914+ if malicious_event :
915+ self .malicious_events .append (malicious_event )
916+
857917 def _parse_storyline_time (self , time_str : str ) -> datetime :
858918 """Parse storyline event time to absolute datetime.
859919
@@ -1773,8 +1833,9 @@ def _generate_system_traffic(self, current_hour: datetime) -> None:
17731833 'message' : action ,
17741834 })
17751835 elif source_roll < 0.80 :
1776- # sshd — disconnect/keepalive messages
1777- ip = f'10.0.{ rng .randint (0 ,10 )} .{ rng .randint (1 ,254 )} '
1836+ # sshd — disconnect/keepalive messages (use scenario system IPs)
1837+ other_ips = [s .ip for s in self .scenario .environment .systems if s .ip != system .ip ]
1838+ ip = rng .choice (other_ips ) if other_ips else system .ip
17781839 port = rng .randint (49152 , 65535 )
17791840 msgs = [
17801841 f'Received disconnect from { ip } port { port } :11: disconnected by user' ,
0 commit comments