From 2548ba1df5973c4db1905a44f2447f5f5c82eb5a Mon Sep 17 00:00:00 2001 From: hub-gif <2487812171@qq.com> Date: Tue, 14 Apr 2026 13:11:21 +0800 Subject: [PATCH] fix(pipeline): align models and serializers with pause/resume checkpoint Made-with: Cursor --- .../migrations/0012_job_pause_checkpoint.py | 65 +++++++++++++++++++ backend/pipeline/models.py | 22 +++++++ backend/pipeline/serializers.py | 44 +++++++++++-- backend/pipeline/urls.py | 5 ++ 4 files changed, 132 insertions(+), 4 deletions(-) create mode 100644 backend/pipeline/migrations/0012_job_pause_checkpoint.py diff --git a/backend/pipeline/migrations/0012_job_pause_checkpoint.py b/backend/pipeline/migrations/0012_job_pause_checkpoint.py new file mode 100644 index 0000000..e9248f5 --- /dev/null +++ b/backend/pipeline/migrations/0012_job_pause_checkpoint.py @@ -0,0 +1,65 @@ +# Generated manually for cookie pause / resume checkpoint + +import django.db.models.deletion +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ("pipeline", "0011_job_cancel_and_status"), + ] + + operations = [ + migrations.AddField( + model_name="pipelinejob", + name="resume_from_checkpoint", + field=models.BooleanField(db_index=True, default=False), + ), + migrations.AlterField( + model_name="pipelinejob", + name="status", + field=models.CharField( + choices=[ + ("pending", "待执行"), + ("running", "执行中"), + ("success", "成功"), + ("failed", "失败"), + ("cancelled", "已终止"), + ("paused", "已暂停(待换 Cookie 续跑)"), + ], + db_index=True, + default="pending", + max_length=16, + ), + ), + migrations.CreateModel( + name="PipelineJobCheckpoint", + fields=[ + ( + "id", + models.BigAutoField( + auto_created=True, + primary_key=True, + serialize=False, + verbose_name="ID", + ), + ), + ("phase", models.CharField(db_index=True, max_length=32)), + ("payload", models.JSONField(blank=True, default=dict)), + ("hint_zh", models.TextField(blank=True, default="")), + ("updated_at", models.DateTimeField(auto_now=True)), + ( + "job", + models.OneToOneField( + on_delete=django.db.models.deletion.CASCADE, + related_name="checkpoint_row", + to="pipeline.pipelinejob", + ), + ), + ], + options={ + "ordering": ["-updated_at"], + }, + ), + ] diff --git a/backend/pipeline/models.py b/backend/pipeline/models.py index 52b7796..5be4407 100644 --- a/backend/pipeline/models.py +++ b/backend/pipeline/models.py @@ -7,6 +7,7 @@ class JobStatus(models.TextChoices): SUCCESS = "success", "成功" FAILED = "failed", "失败" CANCELLED = "cancelled", "已终止" + PAUSED = "paused", "已暂停(待换 Cookie 续跑)" class PipelineJob(models.Model): @@ -43,6 +44,7 @@ class PipelineJob(models.Model): db_index=True, ) cancellation_requested = models.BooleanField(default=False, db_index=True) + resume_from_checkpoint = models.BooleanField(default=False, db_index=True) run_dir = models.TextField(blank=True, default="") error_message = models.TextField(blank=True, default="") created_at = models.DateTimeField(auto_now_add=True) @@ -55,6 +57,26 @@ class PipelineJob(models.Model): return f"[{self.platform}] {self.keyword} ({self.status})" +class PipelineJobCheckpoint(models.Model): + """Cookie 暂停续跑等场景的断点元数据(与任务一对一)。""" + + job = models.OneToOneField( + PipelineJob, + on_delete=models.CASCADE, + related_name="checkpoint_row", + ) + phase = models.CharField(max_length=32, db_index=True) + payload = models.JSONField(default=dict, blank=True) + hint_zh = models.TextField(blank=True, default="") + updated_at = models.DateTimeField(auto_now=True) + + class Meta: + ordering = ["-updated_at"] + + def __str__(self) -> str: + return f"checkpoint job={self.job_id} phase={self.phase}" + + class JdProduct(models.Model): """京东 SKU 主档:同一 ``platform`` + ``sku_id`` 唯一,多次抓取时覆盖为最新一行合并表数据。""" diff --git a/backend/pipeline/serializers.py b/backend/pipeline/serializers.py index 050af0d..23afe31 100644 --- a/backend/pipeline/serializers.py +++ b/backend/pipeline/serializers.py @@ -5,7 +5,13 @@ from django.conf import settings from rest_framework import serializers from .cookie_paste import normalize_browser_cookie_paste -from .models import JdProduct, JdProductSnapshot, JobStatus, PipelineJob +from .models import ( + JdProduct, + JdProductSnapshot, + JobStatus, + PipelineJob, + PipelineJobCheckpoint, +) # 与 views._safe_file_for_job 中 mapping 一致,供前端展示「数据源是否就绪」 _REPORT_CONFIG_ALLOWED_KEYS = frozenset( @@ -54,6 +60,7 @@ class PipelineJobSerializer(serializers.ModelSerializer): inline_cookie_used = serializers.SerializerMethodField() analysis_artifacts = serializers.SerializerMethodField() + checkpoint = serializers.SerializerMethodField() class Meta: model = PipelineJob @@ -74,6 +81,8 @@ class PipelineJobSerializer(serializers.ModelSerializer): "report_config", "status", "cancellation_requested", + "resume_from_checkpoint", + "checkpoint", "run_dir", "error_message", "analysis_artifacts", @@ -84,8 +93,10 @@ class PipelineJobSerializer(serializers.ModelSerializer): "id", "inline_cookie_used", "analysis_artifacts", + "checkpoint", "status", "cancellation_requested", + "resume_from_checkpoint", "run_dir", "error_message", "created_at", @@ -96,10 +107,24 @@ class PipelineJobSerializer(serializers.ModelSerializer): def get_inline_cookie_used(self, obj: PipelineJob) -> bool: return bool((obj.cookie_text or "").strip()) + def get_checkpoint(self, obj: PipelineJob) -> dict | None: + try: + c = obj.checkpoint_row + except PipelineJobCheckpoint.DoesNotExist: + return None + return { + "phase": c.phase, + "payload": c.payload, + "hint_zh": c.hint_zh, + "updated_at": c.updated_at, + } + def get_analysis_artifacts(self, obj: PipelineJob) -> dict[str, bool] | None: - if obj.status not in (JobStatus.SUCCESS, JobStatus.CANCELLED) or not ( - obj.run_dir or "" - ).strip(): + if obj.status not in ( + JobStatus.SUCCESS, + JobStatus.CANCELLED, + JobStatus.PAUSED, + ) or not (obj.run_dir or "").strip(): return None try: base = Path(obj.run_dir).expanduser().resolve() @@ -186,6 +211,17 @@ def _jd_data_root() -> Path: return (Path(root) / "data" / "JD").resolve() +class JobResumeRequestSerializer(serializers.Serializer): + """从断点续跑时可选更新 Cookie。""" + + cookie_text = serializers.CharField( + required=False, + allow_blank=True, + default="", + max_length=500_000, + ) + + class CreatePipelineJobSerializer(serializers.Serializer): keyword = serializers.CharField(max_length=256, trim_whitespace=True) platform = serializers.ChoiceField(choices=["jd"], default="jd") diff --git a/backend/pipeline/urls.py b/backend/pipeline/urls.py index 14d09e9..e4d7538 100644 --- a/backend/pipeline/urls.py +++ b/backend/pipeline/urls.py @@ -15,6 +15,11 @@ urlpatterns = [ views.JobCancelView.as_view(), name="job-cancel", ), + path( + "jobs//resume/", + views.JobResumeView.as_view(), + name="job-resume", + ), path("jobs//download/", views.JobDownloadView.as_view(), name="job-download"), path("jobs//preview/", views.JobPreviewView.as_view(), name="job-preview"), path(