1818
1919from croniter import croniter
2020from dbos import DBOS , Queue , SetEnqueueOptions , SetWorkflowAttributes , SetWorkflowID , StepOptions
21- from dbos ._dbos import _get_dbos_instance , _get_or_create_dbos_registry
21+ from dbos ._dbos import _get_or_create_dbos_registry
2222from dbos ._error import (
2323 DBOSAwaitedWorkflowCancelledError ,
24- DBOSQueueDeduplicatedError ,
2524 DBOSWorkflowCancelledError ,
2625)
2726from pydantic import BaseModel , Field , create_model
@@ -941,10 +940,13 @@ async def start(
941940 # DBOS queue deduplication: the slot is claimed atomically at enqueue,
942941 # held while the workflow is enqueued or pending (a parked run keeps
943942 # it), and freed by DBOS itself at the terminal outcome — including
944- # when DBOS gives up on a dead workflow. A duplicate start() hands back
945- # the live run 's id . Subjectless runs are unbounded.
943+ # when DBOS gives up on a dead workflow. On a held slot, return-existing
944+ # hands back the holder 's handle . Subjectless runs are unbounded.
946945 enqueue_options = (
947- SetEnqueueOptions (deduplication_id = f"{ cls .kind } :{ subject .subject_type } :{ subject .id } " )
946+ SetEnqueueOptions (
947+ deduplication_id = f"{ cls .kind } :{ subject .subject_type } :{ subject .id } " ,
948+ duplication_policy = "return-existing" ,
949+ )
948950 if subject
949951 else nullcontext ()
950952 )
@@ -959,30 +961,23 @@ async def start(
959961 "subject_label" : subject .label ,
960962 }
961963 subject_record = subject .identity if subject else None
962- try :
963- with (
964- SetWorkflowID (workflow_id ),
965- SetWorkflowAttributes (attributes ),
966- enqueue_options ,
967- ):
968- await run_queue .enqueue_async (cls ._entry , subject_record , wire )
969- except DBOSQueueDeduplicatedError as duplicate :
970- holder = _get_dbos_instance ()._sys_db .get_deduplicated_workflow (
971- run_queue .name , duplicate .deduplication_id
964+ with (
965+ SetWorkflowID (workflow_id ),
966+ SetWorkflowAttributes (attributes ),
967+ enqueue_options ,
968+ ):
969+ handle = await run_queue .enqueue_async (cls ._entry , subject_record , wire )
970+ if handle .workflow_id == workflow_id :
971+ # The body also creates its row (idempotently) — this one just makes it
972+ # visible before an executor picks the workflow up.
973+ Run .create_row (
974+ _step_engine (), workflow_id = workflow_id , kind = cls .kind , account_id = account_id
972975 )
973- if holder :
974- return holder
975- # The holder reached terminal between the rejection and the lookup —
976- # the slot is free now, so this start goes through.
977- return await cls .start (subject = subject , account_id = account_id , ** input )
978- # The body also creates its row (idempotently) — this one just makes it
979- # visible before an executor picks the workflow up.
980- Run .create_row (
981- _step_engine (), workflow_id = workflow_id , kind = cls .kind , account_id = account_id
982- )
983- if subject :
984- await publish (WorkflowEvent .SCHEDULED , subject = subject .identity , kind = cls .kind )
985- return workflow_id
976+ if subject :
977+ await publish (WorkflowEvent .SCHEDULED , subject = subject .identity , kind = cls .kind )
978+ return workflow_id
979+ # The slot was held — the handle is the subject's live run.
980+ return handle .workflow_id
986981
987982
988983def _wrap_steps (cls : type [Workflow ]) -> None :
0 commit comments