{"record":{"id":"ac2284aa512edf2a","repo":"apache/beam","slug":"already-staging-s","errorCode":null,"errorMessage":"Already staging %s","messagePattern":"Already staging (.+?)","errorType":"validation","errorClass":"ValueError","httpStatus":null,"severity":"error","filePath":"sdks/python/apache_beam/runners/portability/artifact_service.py","lineNumber":112,"sourceCode":"    beam_artifact_api_pb2_grpc.ArtifactStagingServiceServicer):\n  def __init__(\n      self,\n      file_writer: Callable[[str, Optional[str]], tuple[BinaryIO, str]],\n  ):\n    self._lock = threading.Lock()\n    self._jobs_to_stage: dict[\n        str,\n        tuple[dict[Any, list[beam_runner_api_pb2.ArtifactInformation]],\n              threading.Event]] = {}\n    self._file_writer = file_writer\n\n  def register_job(\n      self,\n      staging_token: str,\n      dependency_sets: MutableMapping[\n          Any, list[beam_runner_api_pb2.ArtifactInformation]]):\n    if staging_token in self._jobs_to_stage:\n      raise ValueError('Already staging %s' % staging_token)\n    with self._lock:\n      self._jobs_to_stage[staging_token] = (\n          dict(dependency_sets), threading.Event())\n\n  def resolved_deps(self, staging_token, timeout=None):\n    with self._lock:\n      dependency_sets, event = self._jobs_to_stage[staging_token]\n    try:\n      if not event.wait(timeout):\n        raise concurrent.futures.TimeoutError()\n      return dependency_sets\n    finally:\n      with self._lock:\n        del self._jobs_to_stage[staging_token]\n\n  def ReverseArtifactRetrievalService(self, responses, context=None):\n    staging_token = next(responses).staging_token\n    with self._lock:","sourceCodeStart":94,"sourceCodeEnd":130,"githubUrl":"https://github.com/apache/beam/blob/12126d8942aaf848030c478b4c6a28c6af861c66/sdks/python/apache_beam/runners/portability/artifact_service.py#L94-L130","documentation":"register_job records a staging token so the artifact service can track dependency resolution per job. It raises ValueError if the same staging_token is already registered, preventing duplicate registration and silent overwriting of dependency sets.","triggerScenarios":"Calling register_job (or the stage RPC path that invokes it) twice with the same staging_token — e.g. a client retrying a staging invocation whose first attempt already registered the token, or two jobs sharing an invocation id.","commonSituations":"Client-side retries after network failures mid-staging; reusing an environment/invocation id across runs; test harnesses invoking register_job manually more than once.","solutions":["Generate a unique staging token per staging session; do not reuse invocation ids.","Check 'staging_token in service._jobs_to_stage' before calling register_job.","Catch ValueError and treat it as 'already registered' if idempotent re-staging is intended."],"exampleFix":"// before\nservice.register_job(token, deps)  # second call with same token\n// after\nif token not in service._jobs_to_stage:\n  service.register_job(token, deps)","handlingStrategy":"validation","validationCode":"if staging_token in service._jobs_to_stage:\n    raise RuntimeError(f\"{staging_token} already registered\")","typeGuard":null,"tryCatchPattern":"try:\n    service.register_job(staging_token, deps)\nexcept ValueError:\n    pass  # idempotent re-registration: token already registered","preventionTips":["Generate a unique staging/invocation token per staging session.","Do not blindly retry staging calls without a fresh token.","Use uuid4-based tokens in test harnesses to avoid collisions."],"tags":["python","apache-beam","artifact-staging","duplicate-key"],"backgroundTag":"file-already-exists","analyzedSha":"12126d8942aaf848030c478b4c6a28c6af861c66","analyzedAt":"2026-09-13T01:50:10.254Z","contentChangedAt":"2026-09-13T01:50:10.254Z","schemaVersion":2},"datasetVersion":"2026-09-20T03:17:13.778Z"}