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
2 changes: 1 addition & 1 deletion THIRD_PARTY_SBOM.cyclonedx.json
Original file line number Diff line number Diff line change
Expand Up @@ -113,7 +113,7 @@
},
{
"name": "context-engine:file-sha256",
"value": "third_party/ragflow/deepdoc/parser/docx_parser.py=e84ea01662ce60180e26dde1ca3ec36fc28a9e9d40357f0e9311ce6937a874c5"
"value": "third_party/ragflow/deepdoc/parser/docx_parser.py=b34d6b2cef003e50f23a0b2c587e451476222b44d18a0b48acc750d79abb01b0"
},
{
"name": "context-engine:file-sha256",
Expand Down
45 changes: 7 additions & 38 deletions tests/integration/test_file_import_tracer.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,6 @@
import json
from collections.abc import Iterator
from contextlib import contextmanager
from dataclasses import replace
from datetime import UTC, datetime, timedelta
from hashlib import sha256
from pathlib import Path
Expand Down Expand Up @@ -516,42 +515,6 @@ def _job_state(
migration_engine.dispose()


def _expire_redeemed_lease(
migration_configuration: DatabaseConfiguration,
scenario: _FileImportScenario,
claims: WorkerLeaseClaims,
) -> WorkerLeaseClaims:
migration_engine = create_database_engine(migration_configuration)
try:
with migration_engine.begin() as connection:
row = connection.execute(
text(
"""
UPDATE file_import_job
SET lease_issued_at = date_trunc('second', statement_timestamp())
- interval '20 minutes',
lease_redeemed_at = date_trunc('second', statement_timestamp())
- interval '19 minutes',
lease_expires_at = date_trunc('second', statement_timestamp())
- interval '10 minutes'
WHERE organization_id = :org AND job_id = :job_id
RETURNING lease_issued_at, lease_expires_at
"""
),
{
"org": scenario.organization_id,
"job_id": scenario.prepared.job_id,
},
).one()
finally:
migration_engine.dispose()
return replace(
claims,
issued_at=row.lease_issued_at,
expires_at=row.lease_expires_at,
)


def _scenario_effect_counts(
migration_configuration: DatabaseConfiguration,
scenario: _FileImportScenario,
Expand Down Expand Up @@ -2200,10 +2163,16 @@ def test_expired_redeemed_lease_cannot_publish_or_record_failure(
tmp_path,
migration_configuration,
guarded_control_engine,
lease_ttl_seconds=1,
)
claims = _scenario_claims(scenario)
assert _redeem_direct(guarded_worker_engine, claims) is not None
claims = _expire_redeemed_lease(migration_configuration, scenario, claims)
migration_engine = create_database_engine(migration_configuration)
try:
with migration_engine.connect() as connection:
connection.execute(text("SELECT pg_sleep(1.1)"))
finally:
migration_engine.dispose()

assert (
_publish_direct(
Expand Down
Loading
Loading