(state: BridgeState, interval: float)
| 3122 | _grp.run_dirs.pop(agent.id, None) |
| 3123 | _grp.completed_s7.discard(agent.id) |
| 3124 | remaining = [state.agents.get(a) for a in _grp.agent_ids if state.agents.get(a)] |
| 3125 | waiting = [a for a in remaining if a.status == "waiting_discussion"] |
| 3126 | if waiting and len(_grp.agent_ids) < 2: |
| 3127 | sole = waiting[0] |
| 3128 | all_messages.append(msg_log( |
| 3129 | sole, |
| 3130 | f"伙伴 agent 失败,跳过讨论 → 直接进入 S8 假设生成", |
| 3131 | "warning", DISCUSSION_STAGE, |
| 3132 | )) |
| 3133 | sole.stage_progress[DISCUSSION_STAGE] = "skipped" |
| 3134 | all_messages.append(msg_stage_update(sole.id, DISCUSSION_STAGE, "skipped")) |
| 3135 | all_messages.extend(_launch_s8_for_agent(state, sole, _grp)) |
| 3136 | _grp.status = "done" |
| 3137 | elif not remaining: |
| 3138 | del state.discussion_groups[_disc_key] |
| 3139 | _reset_agent_idle(agent) |
| 3140 | all_messages.append(msg_agent_update(agent)) |
| 3141 | |
| 3142 | # Poll active discussions |
| 3143 | for group in list(state.discussion_groups.values()): |
| 3144 | all_messages.extend(_poll_discussion(state, group)) |
| 3145 | |
| 3146 | # Schedule idle agents |
| 3147 | sched_msgs = schedule_idle_agents(state) |
| 3148 | all_messages.extend(sched_msgs) |
| 3149 | |
| 3150 | # Periodically broadcast project list (every ~10 poll cycles) |
| 3151 | if not hasattr(state, '_project_list_counter'): |
| 3152 | state._project_list_counter = 0 # type: ignore[attr-defined] |
| 3153 | state._project_list_counter += 1 # type: ignore[attr-defined] |
| 3154 | if state._project_list_counter >= 10: # type: ignore[attr-defined] |
| 3155 | state._project_list_counter = 0 # type: ignore[attr-defined] |
| 3156 | all_messages.append(msg_project_list(list_all_projects(state))) |
| 3157 | |
| 3158 | await broadcast(state, all_messages) |
| 3159 | |
| 3160 | |
| 3161 | # ── Startup ───────────────────────────────────────────────────────────────── |
| 3162 | |
| 3163 | async def main(args: argparse.Namespace): |
| 3164 | state = BridgeState( |
| 3165 | python_path=args.python, |
| 3166 | agent_package_dir=args.agent_dir, |
| 3167 | runs_base_dir=args.runs_dir, |
| 3168 | gpu_allocator=GpuAllocator(args.total_gpus, args.gpus_per_project), |
| 3169 | auto_loop=args.auto_loop, |
| 3170 | discussion_mode=args.discussion_mode, |
| 3171 | discussion_rounds=args.discussion_rounds, |
| 3172 | discussion_models=[m.strip() for m in args.discussion_models.split(",") if m.strip()], |
| 3173 | idea_factory_topic=args.idea_topic, |
| 3174 | idea_factory_config=args.idea_config, |
| 3175 | idea_factory_remaining=args.idea_count, |
| 3176 | ) |
| 3177 | |
| 3178 | # Initialize shared results registry |
| 3179 | _shared_results_path = Path(state.runs_base_dir).parent / "shared_results" |
| 3180 | try: |
| 3181 | from result_registry import ResultRegistry |
no test coverage detected