Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
44 changes: 41 additions & 3 deletions packit_service/worker/handlers/distgit.py
Original file line number Diff line number Diff line change
Expand Up @@ -238,6 +238,8 @@ def __init__(
event: dict,
celery_task: Task,
sync_release_run_id: Optional[int] = None,
retry_tag: Optional[str] = None,
retry_version: Optional[str] = None,
):
super().__init__(
package_config=package_config,
Expand All @@ -247,6 +249,8 @@ def __init__(
)
self._project_url = self.data.project_url
self._sync_release_run_id = sync_release_run_id
self._retry_tag = retry_tag
self._retry_version = retry_version
self.helper: Optional[SyncReleaseHelper] = None

@property
Expand Down Expand Up @@ -316,6 +320,8 @@ def sync_branch(
# is not retried also automatically
kargs = self.celery_task.task.request.kwargs.copy()
kargs["sync_release_run_id"] = model.id
kargs["retry_tag"] = tag
kargs["retry_version"] = version
# https://docs.celeryq.dev/en/stable/userguide/tasks.html#retrying
# https://docs.celeryq.dev/en/stable/reference/celery.app.task.html#celery.app.task.Task.retry
self.celery_task.task.retry(
Expand Down Expand Up @@ -344,8 +350,30 @@ def sync_branch(

return downstream_pr, additional_prs

def _get_or_create_sync_release_run(self, project_event_model=None) -> SyncReleaseModel:
if self._sync_release_run_id is not None:
def _get_or_create_sync_release_run(
self,
project_event_model: Optional[ProjectEventModel] = None,
tag: Optional[str] = None,
version: Optional[str] = None,
) -> SyncReleaseModel:
Comment thread
nforro marked this conversation as resolved.
"""Get the existing sync release run model if retrying, or create a new one.

On retry, the existing model is reused only when the tag and version
match the retry parameters, so that other versions in a multi-version
run get their own models.

Args:
project_event_model: The project event model associated with the run.
tag: The tag of the release being synced.
version: The version of the release being synced.

Returns:
The SyncReleaseModel representing the run.
"""
if self._sync_release_run_id is not None and (
(self._retry_tag is None and self._retry_version is None)
or (tag == self._retry_tag and version == self._retry_version)
):
Comment thread
nforro marked this conversation as resolved.
return SyncReleaseModel.get_by_id(self._sync_release_run_id)

sync_release_model, _ = SyncReleaseModel.create_with_new_run(
Expand Down Expand Up @@ -605,7 +633,9 @@ def _run_for_release(
if tag and self.data.event_type == anitya.NewHotness.event_type()
else None
)
sync_release_run_model = self._get_or_create_sync_release_run(release_event)
sync_release_run_model = self._get_or_create_sync_release_run(
release_event, tag=tag, version=version
)
branches_to_run = [target.branch for target in sync_release_run_model.sync_release_targets]
logger.debug(
f"Branches to run {self.job_config.type} "
Expand Down Expand Up @@ -720,13 +750,17 @@ def __init__(
event: dict,
celery_task: Task,
sync_release_run_id: Optional[int] = None,
retry_tag: Optional[str] = None,
retry_version: Optional[str] = None,
):
super().__init__(
package_config=package_config,
job_config=job_config,
event=event,
celery_task=celery_task,
sync_release_run_id=sync_release_run_id,
retry_tag=retry_tag,
retry_version=retry_version,
)

@staticmethod
Expand Down Expand Up @@ -782,13 +816,17 @@ def __init__(
event: dict,
celery_task: Task,
sync_release_run_id: Optional[int] = None,
retry_tag: Optional[str] = None,
retry_version: Optional[str] = None,
):
super().__init__(
package_config=package_config,
job_config=job_config,
event=event,
celery_task=celery_task,
sync_release_run_id=sync_release_run_id,
retry_tag=retry_tag,
retry_version=retry_version,
)
if self.data.event_type in (pagure.pr.Comment.event_type(),):
# use upstream project URL when retriggering from dist-git PR
Expand Down
8 changes: 8 additions & 0 deletions packit_service/worker/tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -469,12 +469,16 @@ def run_propose_downstream_handler(
package_config: dict,
job_config: dict,
sync_release_run_id: Optional[int] = None,
retry_tag: Optional[str] = None,
retry_version: Optional[str] = None,
):
handler = ProposeDownstreamHandler(
package_config=load_package_config(package_config),
job_config=load_job_config(job_config),
event=event,
sync_release_run_id=sync_release_run_id,
retry_tag=retry_tag,
retry_version=retry_version,
celery_task=self,
)
return get_handlers_task_results(handler.run_job(), event)
Expand All @@ -494,12 +498,16 @@ def run_pull_from_upstream_handler(
package_config: dict,
job_config: dict,
sync_release_run_id: Optional[int] = None,
retry_tag: Optional[str] = None,
retry_version: Optional[str] = None,
):
handler = PullFromUpstreamHandler(
package_config=load_package_config(package_config),
job_config=load_job_config(job_config),
event=event,
sync_release_run_id=sync_release_run_id,
retry_tag=retry_tag,
retry_version=retry_version,
celery_task=self,
)
return get_handlers_task_results(handler.run_job(), event)
Expand Down
203 changes: 203 additions & 0 deletions tests/integration/test_new_hotness_update.py
Original file line number Diff line number Diff line change
Expand Up @@ -579,6 +579,209 @@ class AnityaTestProjectModel(AnityaProjectModel):
assert first_dict_value(results["job"])["success"]


def test_retry_pull_from_upstream_multi_version(new_hotness_update):
"""On retry with multiple versions, only the failed version should reuse its
SyncReleaseModel. Other versions must create their own models."""
new_hotness_update["trigger"]["msg"]["message"]["upstream_versions"] = ["7.0.3", "7.0.4"]

class AnityaTestProjectModel(AnityaProjectModel):
pass

db_project_object = flexmock(
id=12,
project_event_model_type=ProjectEventModelType.anitya_multiple_versions,
job_config_trigger_type=JobConfigTriggerType.release,
project=AnityaTestProjectModel(),
)
project_event = (
flexmock(type=ProjectEventModelType.anitya_multiple_versions, event_id=12)
.should_receive("get_project_event_object")
.and_return(db_project_object)
.mock()
)
run_model = flexmock(PipelineModel)
flexmock(ProjectEventModel).should_receive("get_or_create").with_args(
type=ProjectEventModelType.anitya_multiple_versions,
event_id=12,
commit_sha=None,
).and_return(project_event)
flexmock(AnityaMultipleVersionsModel).should_receive("get_or_create").with_args(
versions=["7.0.3", "7.0.4"],
project_name="redis",
project_id=4181,
package="redis",
).and_return(db_project_object)

# The retried model (for version 7.0.4 which failed download previously)
retried_target = flexmock(status="retry", id=1234, branch="main")
retried_model = flexmock(id=123, sync_release_targets=[retried_target])
flexmock(SyncReleaseModel).should_receive("get_by_id").with_args(123).and_return(
retried_model,
).once()

# A new model should be created for version 7.0.3 (never attempted before)
new_target = flexmock(status="queued", id=1235, branch="main")
new_model = flexmock(id=124, sync_release_targets=[])
flexmock(SyncReleaseModel).should_receive("create_with_new_run").with_args(
status=SyncReleaseStatus.running,
project_event_model=project_event,
job_type=SyncReleaseJobType.pull_from_upstream,
package_name="redis",
).and_return(new_model, run_model).once()
flexmock(SyncReleaseTargetModel).should_receive("create").with_args(
status=SyncReleaseTargetStatus.queued,
branch="main",
).and_return(new_target).once()

flexmock(SyncReleasePullRequestModel).should_receive("get_or_create").with_args(
pr_id=21,
namespace="downstream-namespace",
repo_name="downstream-repo",
project_url="https://src.fedoraproject.org/rpms/downstream-repo",
target_branch=str,
url=str,
).and_return(flexmock(sync_release_targets=[flexmock()]))

packit_yaml = (
"{'specfile_path': 'hello-world.spec', "
"jobs: [{trigger: release, job: pull_from_upstream, metadata: {targets:[]}}]}"
)
flexmock(Github, get_repo=lambda full_name_or_id: None)
distgit_project = flexmock(
get_files=lambda ref, recursive: [".packit.yaml"],
get_file_content=lambda path, ref, headers: packit_yaml,
full_repo_name="rpms/redis",
repo="redis",
namespace="rpms",
is_private=lambda: False,
default_branch="main",
)

lp = flexmock(LocalProject, refresh_the_arguments=lambda: None)
flexmock(LocalProjectBuilder, _refresh_the_state=lambda *args: lp)
lp.working_dir = ""
lp.git_project = distgit_project
flexmock(DistGit).should_receive("local_project").and_return(lp)

flexmock(Allowlist, check_and_report=True)

flexmock(
packit_service.worker.handlers.distgit,
).should_receive("get_monitoring_metadata").and_return(
flexmock(all_versions=True),
)

service_config = ServiceConfig().get_service_config()
flexmock(service_config).should_receive("get_project").with_args(
"https://src.fedoraproject.org/rpms/redis",
required=False,
).and_return(distgit_project)
flexmock(service_config).should_receive("get_project").with_args(
"https://src.fedoraproject.org/rpms/redis",
).and_return(distgit_project)

target_project = (
flexmock(namespace="downstream-namespace", repo="downstream-repo")
.should_receive("get_web_url")
.and_return("https://src.fedoraproject.org/rpms/downstream-repo")
.mock()
)
pr = (
flexmock(
id=21,
url="some_url",
target_project=target_project,
description="some-title",
)
.should_receive("comment")
.mock()
)
# 7.0.4 (retried version) succeeds this time
flexmock(PackitAPI).should_receive("sync_release").with_args(
dist_git_branch="main",
versions=["7.0.4"],
create_pr=True,
local_pr_branch_suffix="update-pull_from_upstream",
use_downstream_specfile=True,
add_pr_instructions=True,
resolved_bugs=["rhbz#2106196"],
release_monitoring_project_id=4181,
sync_acls=True,
pr_description_footer=DistgitAnnouncement.get_announcement(),
add_new_sources=True,
fast_forward_merge_branches=set(),
).and_return((pr, {})).once().ordered()
# 7.0.3 (new version) also succeeds
flexmock(PackitAPI).should_receive("sync_release").with_args(
dist_git_branch="main",
versions=["7.0.3"],
create_pr=True,
local_pr_branch_suffix="update-pull_from_upstream",
use_downstream_specfile=True,
add_pr_instructions=True,
resolved_bugs=["rhbz#2106196"],
release_monitoring_project_id=4181,
sync_acls=True,
pr_description_footer=DistgitAnnouncement.get_announcement(),
add_new_sources=True,
fast_forward_merge_branches=set(),
).and_return((pr, {})).once().ordered()
flexmock(PackitAPI).should_receive("clean")

for target, model in [
(retried_target, retried_model),
(new_target, new_model),
]:
flexmock(target).should_receive("set_status").with_args(
status=SyncReleaseTargetStatus.running,
).once()
flexmock(target).should_receive("set_downstream_pr_url").with_args(
downstream_pr_url="some_url",
)
flexmock(target).should_receive("set_downstream_prs").with_args(
downstream_prs=list,
).once()
flexmock(target).should_receive("set_status").with_args(
status=SyncReleaseTargetStatus.submitted,
).once()
flexmock(target).should_receive("set_start_time").once()
flexmock(target).should_receive("set_finished_time").once()
flexmock(target).should_receive("set_logs").once()
flexmock(model).should_receive("set_status").with_args(
status=SyncReleaseStatus.finished,
).once()
model.should_receive("get_package_name").and_return(None)

flexmock(IsRunConditionSatisfied).should_receive("pre_check").and_return(True)

flexmock(AddReleaseEventToDb).should_receive("db_project_object").and_return(
flexmock(
job_config_trigger_type=JobConfigTriggerType.release,
id=123,
project_event_model_type=ProjectEventModelType.release,
),
)
flexmock(group).should_receive("apply_async").once()
flexmock(Pushgateway).should_receive("push").times(2).and_return()
flexmock(shutil).should_receive("rmtree").with_args("")

processing_results = SteveJobs().process_message(new_hotness_update)
event_dict, _, job_config, package_config = get_parameters_from_results(
processing_results,
)
assert json.dumps(event_dict)

# Simulate retry: pass sync_release_run_id and retry_version for 7.0.4
results = run_pull_from_upstream_handler(
package_config=package_config,
event=event_dict,
job_config=job_config,
sync_release_run_id=123,
retry_version="7.0.4",
)
assert first_dict_value(results["job"])["success"]


def test_new_hotness_update_non_git(new_hotness_update, sync_release_model_non_git):
model = flexmock(status="queued", id=1234, branch="main")
flexmock(SyncReleaseTargetModel).should_receive("create").with_args(
Expand Down
Loading