Skip to content

Commit adbd4f3

Browse files
authored
[ENG-11727] Corrupted file versions are unpreviewable, undownloadable, and undeletable (#11820)
## Ticket https://openscience.atlassian.net/browse/ENG-11727 ## Purpose A pattern of file version corruption has emerged across multiple projects. After a new version is uploaded, the file becomes unpreviewable, undownloadable, and impossible to delete — by the user and by support. Uploading additional versions doesn't fix it; metadata updates (date modified) but the file remains unretrievable. The broken versions are stuck and block users from recovering a working file. Looks like what’s happening is if you check https://api.osf.io/v2/files/f42sc/ and look at its file versions https://api.osf.io/v2/files/697ec8ac74e08d03dccd76ea/versions/ we have multiple version 2s of the same file that seem identical but have a sliiightly different date_created, and that's causing problems with the system. We need to prevent that from happening and clean the data for the existing issues. ## Changes make create_version to avoid file version duplicates and get_version to return single file version; create management command to dedupe (delete file versions duplicates);
1 parent b5f689e commit adbd4f3

6 files changed

Lines changed: 277 additions & 16 deletions

File tree

addons/osfstorage/models.py

Lines changed: 14 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@
1414
from osf.models import AbstractNode
1515
from osf.models.files import File, FileVersion, Folder, TrashedFileNode, BaseFileNode, BaseFileNodeManager
1616
from osf.utils import permissions
17+
from osf.utils.requests import check_select_for_update
1718
from website.files import exceptions
1819
from website.files import utils as files_utils
1920
from website.util import api_url_for
@@ -327,6 +328,12 @@ def update_region_from_latest_version(self, destination_parent):
327328
most_recent_fileversion.save()
328329

329330
def create_version(self, creator, location, metadata=None):
331+
if check_select_for_update():
332+
# Lock the file row for the duration of the request's transaction to avoid
333+
# concurrent/retried requests for the same file to read the same version
334+
# count and insert duplicate identifiers
335+
self.__class__.objects.select_for_update().get(pk=self.pk)
336+
330337
latest_version = self.get_version()
331338
version = FileVersion(identifier=self.versions.count() + 1, creator=creator, location=location)
332339

@@ -354,12 +361,13 @@ def get_version(self, version=None, required=False):
354361
return self.versions.first()
355362
return None
356363

357-
try:
358-
return self.versions.get(identifier=version)
359-
except FileVersion.DoesNotExist:
360-
if required:
361-
raise exceptions.VersionNotFoundError(version)
362-
return None
364+
# .filter().first() better than .get(): some files have more
365+
# than one FileVersion sharing the same identifier, which would
366+
# otherwise raise MultipleObjectsReturned here instead of retrieving a version
367+
result = self.versions.filter(identifier=version).order_by('created').first()
368+
if result is None and required:
369+
raise exceptions.VersionNotFoundError(version)
370+
return result
363371

364372
def add_tag_log(self, action, tag, auth):
365373
if isinstance(self.target, Loggable):

addons/osfstorage/tests/test_models.py

Lines changed: 39 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,8 @@
33

44
import pytest
55
import pytz
6+
from django.db import connection, transaction
7+
from django.test.utils import CaptureQueriesContext
68
from django.utils import timezone
79
from importlib import import_module
810
from django.conf import settings as django_conf_settings
@@ -195,9 +197,43 @@ def test_download_count_file(self):
195197
assert child.get_download_count(1) == 1
196198
assert child.get_download_count(2) == 1
197199

198-
@unittest.skip
199-
def test_create_version(self):
200-
pass
200+
def test_create_version_locks_file_row(self):
201+
202+
file = self.node_settings.get_root().append_file('locked.txt')
203+
204+
with transaction.atomic(), CaptureQueriesContext(connection) as ctx:
205+
file.create_version(
206+
self.user,
207+
{
208+
'service': 'cloud',
209+
settings.WATERBUTLER_RESOURCE: 'osf',
210+
'object': '06d80e',
211+
}, {
212+
'size': 1234,
213+
'contentType': 'text/plain'
214+
})
215+
216+
for_update_sql = connection.ops.for_update_sql()
217+
assert any(for_update_sql in query['sql'] for query in ctx.captured_queries)
218+
219+
@mock.patch('osf.utils.requests.settings.SELECT_FOR_UPDATE_ENABLED', False)
220+
def test_create_version_does_not_lock_file_row_when_disabled(self):
221+
file = self.node_settings.get_root().append_file('unlocked.txt')
222+
223+
with transaction.atomic(), CaptureQueriesContext(connection) as ctx:
224+
file.create_version(
225+
self.user,
226+
{
227+
'service': 'cloud',
228+
settings.WATERBUTLER_RESOURCE: 'osf',
229+
'object': '06d80e',
230+
}, {
231+
'size': 1234,
232+
'contentType': 'text/plain'
233+
})
234+
235+
for_update_sql = connection.ops.for_update_sql()
236+
assert not any(for_update_sql in query['sql'] for query in ctx.captured_queries)
201237

202238
def test_delete_folder(self):
203239
parent = self.node_settings.get_root().append_folder('Test')
Lines changed: 123 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,123 @@
1+
import logging
2+
3+
from django.core.management.base import BaseCommand
4+
from django.db import transaction
5+
from django.db.models import Count
6+
7+
from osf.models.files import BaseFileNode, BaseFileVersionsThrough
8+
from osf.utils.requests import check_select_for_update
9+
10+
logger = logging.getLogger(__name__)
11+
12+
13+
def find_duplicate_groups():
14+
"""
15+
Yields {'basefilenode_id', 'fileversion__identifier', 'count'} for every
16+
(file, identifier) pair that has more than one linked FileVersion.
17+
"""
18+
return (
19+
BaseFileVersionsThrough.objects
20+
.values('basefilenode_id', 'fileversion__identifier')
21+
.annotate(count=Count('id'))
22+
.filter(count__gt=1)
23+
.order_by('basefilenode_id')
24+
)
25+
26+
27+
def fetch_group_rows(basefilenode_id, identifier):
28+
# `fileversion__id` is a secondary sort key so the choice of keeper is fully
29+
# deterministic even if two duplicates share the exact same `created` timestamp.
30+
return list(
31+
BaseFileVersionsThrough.objects
32+
.filter(basefilenode_id=basefilenode_id, fileversion__identifier=identifier)
33+
.select_related('fileversion')
34+
.order_by('fileversion__created', 'fileversion__id')
35+
)
36+
37+
38+
def log_group(file_label, identifier, row_to_keep, rows_to_delete, dry_run):
39+
logger.info(
40+
f'{"[DRY-RUN] " if dry_run else ""}file={file_label} identifier={identifier} '
41+
f'keeping {row_to_keep.fileversion._id} (location={row_to_keep.fileversion.location}) '
42+
f'discarding={[(row.fileversion._id, row.fileversion.location) for row in rows_to_delete]}'
43+
)
44+
45+
46+
def resolve_group(basefilenode_id, identifier, file_label, dry_run):
47+
"""
48+
Fetches the current rows for one duplicate group, logs the keep/discard
49+
decision, unless dry_run - deletes the discarded duplicate(s.
50+
"""
51+
through_rows = fetch_group_rows(basefilenode_id, identifier)
52+
if len(through_rows) < 2:
53+
return False
54+
55+
row_to_keep, rows_to_delete = through_rows[0], through_rows[1:]
56+
log_group(file_label, identifier, row_to_keep, rows_to_delete, dry_run=dry_run)
57+
58+
if not dry_run:
59+
for row in rows_to_delete:
60+
row.delete()
61+
62+
return True
63+
64+
65+
def dedupe_file_versions(dry_run=True):
66+
"""
67+
Finds FileVersions that share the same `identifier` for the same file because of
68+
race condition in OsfStorageFile.create_version() and delete duplicate
69+
"""
70+
if dry_run:
71+
logger.info('[DRY-RUN] Data will not be modified.')
72+
73+
groups = list(find_duplicate_groups())
74+
file_labels = dict(
75+
BaseFileNode.objects
76+
.filter(id__in={group['basefilenode_id'] for group in groups})
77+
.values_list('id', '_id')
78+
)
79+
80+
fixed = 0
81+
82+
for group in groups:
83+
basefilenode_id = group['basefilenode_id']
84+
identifier = group['fileversion__identifier']
85+
file_label = file_labels.get(basefilenode_id, basefilenode_id)
86+
87+
if dry_run:
88+
resolved = resolve_group(basefilenode_id, identifier, file_label, dry_run=True)
89+
else:
90+
with transaction.atomic():
91+
if check_select_for_update():
92+
# Lock the file row for the duration of this group's cleanup so
93+
# a concurrent create_version() call for the same file can't
94+
# interleave with the read-then-delete in resolve_group().
95+
BaseFileNode.objects.select_for_update().get(pk=basefilenode_id)
96+
resolved = resolve_group(basefilenode_id, identifier, file_label, dry_run=False)
97+
98+
if resolved:
99+
fixed += 1
100+
101+
logger.info(f'{fixed} duplicate group(s) resolved.')
102+
return fixed
103+
104+
105+
class Command(BaseCommand):
106+
help = """
107+
Finds FileVersions that share the same `identifier` on the same file because of
108+
a race condition in OsfStorageFile.create_version() and delete duplicate
109+
"""
110+
111+
def add_arguments(self, parser):
112+
parser.add_argument(
113+
'--apply',
114+
action='store_false',
115+
dest='dry_run',
116+
default=True,
117+
help='Actually unlink duplicate versions. Without this flag, only reports what would change.',
118+
)
119+
120+
# Management command handler
121+
def handle(self, *args, **options):
122+
dry_run = options.get('dry_run', True)
123+
dedupe_file_versions(dry_run=dry_run)

osf/models/files.py

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,6 @@
66
from django.apps import apps
77
from django.db import models, IntegrityError
88
from django.db.models import Manager
9-
from django.core.exceptions import ObjectDoesNotExist
109
from django.utils import timezone
1110
from django.contrib.contenttypes.models import ContentType
1211
from django.contrib.contenttypes.fields import GenericForeignKey
@@ -289,12 +288,13 @@ def get_version(self, revision, required=False):
289288
:returns: FileVersion or None
290289
:raises: VersionNotFoundError if required is True
291290
"""
292-
try:
293-
return self.versions.get(identifier=revision)
294-
except ObjectDoesNotExist:
295-
if required:
296-
raise VersionNotFoundError(revision)
297-
return None
291+
# .filter().first() better than .get(): some files have more
292+
# than one FileVersion sharing the same identifier, which would otherwise raise
293+
# MultipleObjectsReturned here instead of retrieving a version
294+
version = self.versions.filter(identifier=revision).order_by('created').first()
295+
if version is None and required:
296+
raise VersionNotFoundError(revision)
297+
return version
298298

299299
def generate_waterbutler_url(self, **kwargs):
300300
base_url = None
Lines changed: 73 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,73 @@
1+
import pytest
2+
from django.core.management import call_command
3+
4+
from addons.osfstorage import settings as osfstorage_settings
5+
from addons.osfstorage.tests.factories import FileVersionFactory
6+
from osf.models import BaseFileVersionsThrough, FileVersion
7+
from osf_tests.factories import ProjectFactory
8+
9+
10+
def make_location(obj):
11+
return {
12+
'service': 'cloud',
13+
osfstorage_settings.WATERBUTLER_RESOURCE: 'osf',
14+
'object': obj,
15+
}
16+
17+
18+
@pytest.mark.django_db
19+
class TestDedupeFileVersions:
20+
21+
@pytest.fixture()
22+
def file_node(self):
23+
project = ProjectFactory()
24+
return project.get_addon('osfstorage').get_root().append_file('dupes.txt')
25+
26+
@pytest.fixture()
27+
def add_duplicate_version(self, file_node):
28+
def _add_duplicate_version(identifier, location):
29+
version = FileVersionFactory(identifier=identifier, location=location)
30+
file_node.add_version(version)
31+
return version
32+
return _add_duplicate_version
33+
34+
def test_dry_run_leaves_duplicates_untouched(self, file_node, add_duplicate_version):
35+
add_duplicate_version('2', make_location('object-a'))
36+
add_duplicate_version('2', make_location('object-b'))
37+
38+
call_command('dedupe_file_versions')
39+
40+
assert file_node.versions.filter(identifier='2').count() == 2
41+
42+
def test_apply_unlinks_duplicate_keeping_earliest(self, file_node, add_duplicate_version):
43+
version_to_keep = add_duplicate_version('2', make_location('object-a'))
44+
extra = add_duplicate_version('2', make_location('object-b'))
45+
46+
call_command('dedupe_file_versions', dry_run=False)
47+
48+
remaining = list(file_node.versions.filter(identifier='2'))
49+
assert remaining == [version_to_keep]
50+
assert FileVersion.objects.filter(id=extra.id).exists()
51+
assert not BaseFileVersionsThrough.objects.filter(basefilenode=file_node, fileversion=extra).exists()
52+
53+
def test_leaves_non_duplicate_versions_alone(self, file_node, add_duplicate_version):
54+
add_duplicate_version('1', make_location('object-1'))
55+
add_duplicate_version('2', make_location('object-2'))
56+
57+
call_command('dedupe_file_versions', dry_run=False)
58+
59+
assert file_node.versions.count() == 2
60+
61+
def test_apply_deletes_every_duplicate_but_the_keeper(self, file_node, add_duplicate_version):
62+
version_to_keep = add_duplicate_version('2', make_location('object-a'))
63+
extra_1 = add_duplicate_version('2', make_location('object-b'))
64+
extra_2 = add_duplicate_version('2', make_location('object-c'))
65+
66+
assert BaseFileVersionsThrough.objects.filter(basefilenode=file_node).count() == 3
67+
68+
call_command('dedupe_file_versions', dry_run=False)
69+
70+
assert BaseFileVersionsThrough.objects.filter(basefilenode=file_node).count() == 1
71+
assert list(file_node.versions.all()) == [version_to_keep]
72+
for extra in (extra_1, extra_2):
73+
assert not BaseFileVersionsThrough.objects.filter(basefilenode=file_node, fileversion=extra).exists()

osf_tests/test_files.py

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -198,3 +198,24 @@ def test_file_shared_purged(project, create_test_file):
198198
assert freed == sum(list(test_file_1.versions.values_list('size', flat=True)))
199199
assert test_file_1.purged is not None
200200
assert version_0.purged is not None
201+
202+
203+
def test_get_version_resolves_duplicate_identifiers(project, create_test_file):
204+
# Simulates the historical create_version race that could attach two
205+
# FileVersions with the same identifier to one file. get_version() must
206+
# resolve this instead of raising MultipleObjectsReturned.
207+
test_file = create_test_file(target=project)
208+
first_version = test_file.versions.first()
209+
210+
duplicate_version = FileVersion(
211+
creator=first_version.creator,
212+
identifier=first_version.identifier,
213+
location=dict(first_version.location, object='deadbeef'),
214+
)
215+
duplicate_version.save()
216+
test_file.add_version(duplicate_version)
217+
218+
assert test_file.versions.filter(identifier=first_version.identifier).count() == 2
219+
220+
resolved = test_file.get_version(first_version.identifier, required=True)
221+
assert resolved == first_version

0 commit comments

Comments
 (0)