from statebase import StateBase
sb = StateBase(api_key="your-key")
def orchestrator(session_id, task):
# 1. Decompose
subtasks = llm.decompose(task)
# 2. Record the plan in shared state
sb.sessions.update_state(
session_id=session_id,
state={"plan": subtasks, "results": {}},
reasoning=f"Decomposed task into {len(subtasks)} subtasks"
)
# 3. Dispatch workers (each writes to the same session)
for i, subtask in enumerate(subtasks):
worker_id = dispatch_worker(session_id, i, subtask)
sb.sessions.update_state(
session_id=session_id,
state={f"worker_{i}": worker_id},
reasoning=f"Dispatched worker {i}"
)
# 4. Wait for all workers, merge results
while not all_workers_done(session_id):
time.sleep(5)
results = sb.sessions.get(session_id).state["results"]
return merge(results)