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
78 changes: 74 additions & 4 deletions cds/modules/deposit/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,7 @@
get_tasks_status_grouped_by_task_name,
merge_tasks_status,
)
from ..flows.tasks import ExtractChapterFramesTask
from ..flows.models import FlowMetadata
from ..invenio_deposit.api import Deposit, has_status, preserve
from ..invenio_deposit.utils import mark_as_action
Expand All @@ -76,7 +77,7 @@
)
from ..records.minters import cds_doi_generator, is_local_doi, report_number_minter
from ..records.resolver import record_resolver
from ..records.utils import is_record, lowercase_value
from ..records.utils import is_record, lowercase_value, parse_video_chapters
from ..records.validators import PartialDraft4Validator
from ..records.permissions import is_public
from .errors import DiscardConflict
Expand Down Expand Up @@ -504,7 +505,7 @@ def create(cls, data, id_=None, **kwargs):
data.setdefault("_access", {})
access_update = data["_access"].setdefault("update", [])
try:
if current_user.email not in access_update:
if current_user.email not in access_update:
# Add the current user to the ``_access.update`` list
access_update.append(current_user.email)
except AttributeError:
Expand Down Expand Up @@ -905,11 +906,74 @@ def _publish_edited(self):

return super(Video, self)._publish_edited()

def _has_chapters_changed(self, old_record=None):
"""Check if chapters in description have changed."""
current_description = self.get("description", "")
current_chapters = parse_video_chapters(current_description)

if old_record is None:
# First publish - trigger if chapters exist
return len(current_chapters) > 0

old_description = old_record.get("description", "")
old_chapters = parse_video_chapters(old_description)

# Compare chapter timestamps and titles
if len(current_chapters) != len(old_chapters):
return True

for curr, old in zip(current_chapters, old_chapters):
if curr["seconds"] != old["seconds"] or curr["title"] != old["title"]:
return True

return False

def _trigger_chapter_frame_extraction(self):
"""Trigger chapter frame extraction asynchronously for existing video files."""
try:
# Get the current flow for this deposit
current_flow = FlowMetadata.get_by_deposit(self["_deposit"]["id"])

if current_flow is None:
current_app.logger.warning(
f"No current flow found for video {self.id}. Cannot trigger chapter frame extraction."
)
return

current_app.logger.info(
f"Triggering asynchronous ExtractChapterFramesTask for video {self.id} with flow {current_flow.id}"
)

payload = current_flow.payload.copy()

current_app.logger.info(f"Submitting ExtractChapterFramesTask with payload: {payload}")

ExtractChapterFramesTask().s(**payload).apply_async()

current_app.logger.info(
f"ExtractChapterFramesTask submitted asynchronously for video {self.id}, flow_id: {current_flow.id}"
)
except Exception as e:
current_app.logger.error(
f"Failed to trigger async chapter frame extraction for video {self.id}: {e}"
)
import traceback

current_app.logger.error(f"Traceback: {traceback.format_exc()}")

@mark_as_action
def publish(self, pid=None, id_=None, **kwargs):
"""Publish a video and update the related project."""
# save a copy of the old PID
video_old_id = self["_deposit"]["id"]

# Check if this is a republish and get the old record
old_record = None
try:
Comment thread
zzacharo marked this conversation as resolved.
_, old_record = self.fetch_published()
except KeyError as e: # First publish (no pid key)
pass

try:
self["category"] = self.project["category"]
self["type"] = self.project["type"]
Expand All @@ -930,6 +994,13 @@ def publish(self, pid=None, id_=None, **kwargs):
video_published = super(Video, self).publish(pid=pid, id_=id_, **kwargs)
_, record_new = self.fetch_published()

# Check if chapters have changed and trigger frame extraction
if self._has_chapters_changed(old_record):
current_app.logger.info(
f"Chapters changed for video {self.id}, triggering frame extraction"
)
self._trigger_chapter_frame_extraction()

# update associated project
video_published.project._update_videos(
[video_build_url(video_old_id)],
Expand Down Expand Up @@ -1088,7 +1159,6 @@ def _create_tags(self):
except IndexError:
return


def mint_doi(self):
"""Mint DOI."""
assert self.has_record()
Expand All @@ -1109,7 +1179,7 @@ def mint_doi(self):
status=PIDStatus.RESERVED,
)
return self


project_resolver = Resolver(
pid_type="depid",
Expand Down
2 changes: 2 additions & 0 deletions cds/modules/deposit/receivers.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@
from cds.modules.flows.tasks import (
DownloadTask,
ExtractFramesTask,
ExtractChapterFramesTask,
ExtractMetadataTask,
TranscodeVideoTask,
)
Expand Down Expand Up @@ -87,4 +88,5 @@ def register_celery_class_based_tasks(sender, app=None):
celery.register_task(ExtractMetadataTask())
celery.register_task(DownloadTask())
celery.register_task(ExtractFramesTask())
celery.register_task(ExtractChapterFramesTask())
celery.register_task(TranscodeVideoTask())
2 changes: 2 additions & 0 deletions cds/modules/flows/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@
from .tasks import (
CeleryTask,
DownloadTask,
ExtractChapterFramesTask,
ExtractFramesTask,
ExtractMetadataTask,
TranscodeVideoTask,
Expand Down Expand Up @@ -245,6 +246,7 @@ def _find_celery_task_by_name(name):
ExtractMetadataTask,
ExtractFramesTask,
TranscodeVideoTask,
ExtractChapterFramesTask,
]:
if celery_task.name == name:
return celery_task
Expand Down
Loading