Start an output format manager for this channel if not already running. Only the TS-owning worker starts the manager; non-owners read the shared buffer. Returns True if a manager is active (locally or on another worker).
(self, channel_id, fmt, source_buffer=None)
| 1157 | return self.stream_buffers.get(channel_id) |
| 1158 | |
| 1159 | def ensure_output_format(self, channel_id, fmt, source_buffer=None) -> bool: |
| 1160 | """ |
| 1161 | Start an output format manager for this channel if not already running. |
| 1162 | Only the TS-owning worker starts the manager; non-owners read the shared buffer. |
| 1163 | Returns True if a manager is active (locally or on another worker). |
| 1164 | """ |
| 1165 | if channel_id in self.output_managers and fmt in self.output_managers[channel_id]: |
| 1166 | return True |
| 1167 | |
| 1168 | if not self.redis_client: |
| 1169 | return False |
| 1170 | |
| 1171 | state = self.redis_client.get(RedisKeys.output_state(channel_id, fmt)) |
| 1172 | if state == 'active': |
| 1173 | owner_val = self.redis_client.get(RedisKeys.output_owner(channel_id, fmt)) |
| 1174 | if owner_val and owner_val != self.worker_id: |
| 1175 | logger.info(f"[output:{fmt}] Channel {channel_id}: manager active on another worker") |
| 1176 | return True |
| 1177 | # State says active but we have no local manager - stale state from a dead manager. |
| 1178 | # Fall through to restart if we can. |
| 1179 | logger.warning( |
| 1180 | f"[output:{fmt}] Channel {channel_id}: stale active state detected " |
| 1181 | f"(owner={owner_val}), restarting manager" |
| 1182 | ) |
| 1183 | |
| 1184 | if not self.am_i_owner(channel_id): |
| 1185 | # Ask the TS-owning worker to start the manager, then poll until active. |
| 1186 | logger.info(f"[output:{fmt}] Channel {channel_id}: requesting owner to start manager") |
| 1187 | self.redis_client.publish( |
| 1188 | f"live:events:{channel_id}", |
| 1189 | json.dumps({ |
| 1190 | "event": EventType.ENSURE_OUTPUT_FORMAT, |
| 1191 | "channel_id": channel_id, |
| 1192 | "fmt": fmt, |
| 1193 | "timestamp": time.time(), |
| 1194 | }) |
| 1195 | ) |
| 1196 | deadline = time.time() + 5 |
| 1197 | while time.time() < deadline: |
| 1198 | gevent.sleep(0.1) |
| 1199 | state = self.redis_client.get(RedisKeys.output_state(channel_id, fmt)) |
| 1200 | if state == 'active': |
| 1201 | logger.info(f"[output:{fmt}] Channel {channel_id}: manager started by owner") |
| 1202 | return True |
| 1203 | logger.warning(f"[output:{fmt}] Channel {channel_id}: owner did not start manager within 5s") |
| 1204 | return False |
| 1205 | |
| 1206 | ts_buffer = source_buffer |
| 1207 | if ts_buffer is None: |
| 1208 | _, profile_id = self._parse_output_key(fmt) |
| 1209 | if profile_id is not None: |
| 1210 | ts_buffer = self.profile_buffers.get(channel_id, {}).get(profile_id) |
| 1211 | if ts_buffer is None: |
| 1212 | ts_buffer = self.stream_buffers.get(channel_id) |
| 1213 | if not ts_buffer: |
| 1214 | logger.error(f"[output:{fmt}] Channel {channel_id}: no TS buffer, cannot start manager") |
| 1215 | return False |
| 1216 |
no test coverage detected