• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

uc-cdis / indexd / 30648324736

31 Jul 2026 04:44PM UTC coverage: 87.195% (-0.2%) from 87.405%
30648324736

Pull #442

github

web-flow
Merge 2b6931ec7 into f2275268f
Pull Request #442: Refactor for fastapi

3180 of 3647 relevant lines covered (87.19%)

0.87 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

55.07
indexd/utils.py
1
import re
1✔
2
from urllib.parse import urlparse
1✔
3
import os
1✔
4
import requests
1✔
5
from cdislogging import get_logger
1✔
6
from sqlalchemy import create_engine
1✔
7

8
logger = get_logger(__name__)
1✔
9

10

11
def hint_match(record, hints):
1✔
12
    for hint in hints:
×
13
        if re.match(hint, record):
×
14
            return True
×
15
    return False
×
16

17

18
def try_drop_test_data(
19
    user, database, root_user="postgres", host=""
20
):  # pragma: no cover
21
    engine = create_engine(
22
        "postgresql://{user}@{host}/postgres".format(user=root_user, host=host)
23
    )
24

25
    conn = engine.connect()
26
    conn.execute("commit")
27

28
    try:
29
        create_stmt = 'DROP DATABASE "{database}"'.format(database=database)
30
        conn.execute(create_stmt)
31
    except Exception:
32
        logger.warning("Unable to drop test data:")
33

34
    conn.close()
35

36

37
def setup_database(
38
    user,
39
    password,
40
    database,
41
    root_user="postgres",
42
    host="",
43
    no_drop=False,
44
    no_user=False,
45
):  # pragma: no cover
46
    """
47
    setup the user and database
48
    """
49

50
    if not no_drop:
51
        try_drop_test_data(user, database)
52

53
    engine = create_engine(
54
        "postgresql://{user}@{host}/postgres".format(user=root_user, host=host)
55
    )
56
    conn = engine.connect()
57
    conn.execute("commit")
58

59
    create_stmt = 'CREATE DATABASE "{database}"'.format(database=database)
60
    try:
61
        conn.execute(create_stmt)
62
    except Exception:
63
        logger.warning("Unable to create database")
64

65
    if not no_user:
66
        try:
67
            user_stmt = "CREATE USER {user} WITH PASSWORD '{password}'".format(
68
                user=user, password=password
69
            )
70
            conn.execute(user_stmt)
71

72
            perm_stmt = (
73
                "GRANT ALL PRIVILEGES ON DATABASE {database} to {password}"
74
                "".format(database=database, password=password)
75
            )
76
            conn.execute(perm_stmt)
77
            conn.execute("commit")
78
        except Exception:
79
            logger.warning("Unable to add user")
80
    conn.close()
81

82

83
def create_tables(host, user, password, database):  # pragma: no cover
84
    """
85
    create tables
86
    """
87
    engine = create_engine(
88
        "postgresql://{user}:{pwd}@{host}/{db}".format(
89
            user=user, host=host, pwd=password, db=database
90
        )
91
    )
92
    conn = engine.connect()
93

94
    create_index_record_stm = "CREATE TABLE index_record (\
95
        did VARCHAR NOT NULL, rev VARCHAR, form VARCHAR, size BIGINT, PRIMARY KEY (did) )"
96
    create_record_hash_stm = "CREATE TABLE index_record_hash (\
97
        did VARCHAR NOT NULL, hash_type VARCHAR NOT NULL, hash_value VARCHAR, \
98
        PRIMARY KEY (did, hash_type), FOREIGN KEY(did) REFERENCES index_record (did))"
99
    create_record_url_stm = "CREATE TABLE index_record_url( \
100
        did VARCHAR NOT NULL, url VARCHAR NOT NULL, PRIMARY KEY (did, url),\
101
        FOREIGN KEY(did) REFERENCES index_record (did) )"
102
    create_index_schema_version_stm = "CREATE TABLE index_schema_version (\
103
        version INT)"
104
    create_drs_bundle_record = "CREATE TABLE drs_bundle_record (\
105
        bundle_id VARCHAR NOT NULL, name VARCHAR, created_time DATETIME, updated_time DATETIME,\
106
        checksum VARCHAR, size BIGINT, bundle_data TEXT, description TEXT, version VARCHAR, aliases VARCHAR, PRIMARY KEY(bundle_id)"
107
    try:
108
        conn.execute(create_index_record_stm)
109
        conn.execute(create_record_hash_stm)
110
        conn.execute(create_record_url_stm)
111
        conn.execute(create_index_schema_version_stm)
112
        conn.execute(create_drs_bundle_record)
113
    except Exception:
114
        logger.warning("Unable to create table")
115
        raise
116
    finally:
117
        conn.close()
118

119

120
def check_engine_for_migrate(engine):
1✔
121
    """
122
    check if a db engine support database migration
123

124
    Args:
125
        engine (sqlalchemy.engine.base.Engine): a sqlalchemy engine
126

127
    Return:
128
        bool: whether the engine support migration
129
    """
130
    return engine.dialect.supports_alter
×
131

132

133
def init_schema_version(driver, model, version):
1✔
134
    """
135
    initialize schema table with a initialized singleton of version
136

137
    Args:
138
        driver (object): an alias or index driver instance
139
        model (sqlalchemy.ext.declarative.api.Base): the version table model
140

141
    Return:
142
        version (int): current version number in database
143
    """
144
    with driver.session as s:
×
145
        schema_version = s.query(model).first()
×
146
        if not schema_version:
×
147
            schema_version = model(version=version)
×
148
            s.add(schema_version)
×
149
        version = schema_version.version
×
150
    return version
×
151

152

153
def migrate_database(driver, migrate_functions, current_schema_version, model):
1✔
154
    """
155
    This migration logic is DEPRECATED. It is still supported for backwards compatibility,
156
    but any new migration should be added using Alembic.
157

158
    migrate current database to match the schema version provided in
159
    current schema
160

161
    Args:
162
        driver (object): an alias or index driver instance
163
        migrate_functions (list): a list of migration functions
164
        curent_schema_version (int): version of current schema in code
165
        model (sqlalchemy.ext.declarative.api.Base): the version table model
166

167
    Return:
168
        None
169
    """
170
    db_schema_version = init_schema_version(driver, model, 0)
×
171

172
    need_migrate = (current_schema_version - db_schema_version) > 0
×
173

174
    if not check_engine_for_migrate(driver.engine) and need_migrate:
×
175
        driver.logger.error(
×
176
            "Engine {} does not support alter, skip migration".format(
177
                driver.engine.dialect.name
178
            )
179
        )
180
        return
×
181

182
    for f in migrate_functions[db_schema_version:current_schema_version]:
×
183
        with driver.session as s:
×
184
            schema_version = s.query(model).first()
×
185
            driver.logger.info(
×
186
                "migrating {} schema to {}".format(
187
                    driver.__class__.__name__, schema_version.version
188
                )
189
            )
190

191
            f(engine=driver.engine, session=s)
×
192
            schema_version.version += 1
×
193
            s.add(schema_version)
×
194

195

196
def reverse_url(url):
1✔
197
    """
198
    Reverse the domain name for drs service-info IDs
199
    Args:
200
        url (str): url of the domain
201
        example: drs.example.org
202

203
    returns:
204
        id (str): DRS service-info ID
205
        example: org.example.drs
206
    """
207
    parsed_url = urlparse(url)
1✔
208
    if parsed_url.scheme in ["http", "https"]:
1✔
209
        url = parsed_url.hostname
×
210
    segments = url.split(".")
1✔
211
    reversed_segments = reversed(segments)
1✔
212
    res = ".".join(reversed_segments)
1✔
213
    return res
1✔
214

215

216
FENCE_SERVICE = os.environ.get("FENCE_SERVICE_URL", "http://fence-service")
1✔
217

218

219
def lookup_bucket_region(bucket_name, bucket_regions):
1✔
220
    """
221
    Resolve a bucket name to a region.
222
    Exact match first, then simple prefix fallback.
223
    """
224
    if bucket_name in bucket_regions:
1✔
225
        return bucket_regions[bucket_name]
1✔
226

227
    # remove regexp for prefix matching to remove snyk vulnerability
228
    for pattern, region in bucket_regions.items():
1✔
229
        if pattern.endswith(".*") and bucket_name.startswith(pattern[:-2]):
1✔
230
            return region
1✔
231

232
    return ""
1✔
233

234

235
def get_bucket_regions(app):
1✔
236
    cached = getattr(app, "cache", None)
1✔
237
    if cached:
1✔
238
        hit = cached.get("bucket_regions")
×
239
        if hit:
×
240
            return hit
×
241

242
    url = f"{FENCE_SERVICE}/data/buckets"
1✔
243
    data = {}
1✔
244

245
    try:
1✔
246
        resp = requests.get(url)
1✔
247
        resp.raise_for_status()
×
248
        data = resp.json().get("S3_BUCKETS") or {}
×
249
    except Exception as e:
1✔
250
        logger.warning(f"Failed to fetch bucket regions from Fence: {e}")
1✔
251

252
    regions = {k: v.get("region", "") for k, v in data.items()}
1✔
253

254
    if cached:
1✔
255
        cached.set("bucket_regions", regions)
×
256

257
    return regions
1✔
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc