tighten run-batch manifest replay validation
This commit is contained in:
parent
f0826abb52
commit
68bbf94577
4 changed files with 45 additions and 8 deletions
|
|
@ -307,16 +307,13 @@ class IngesterApp:
|
||||||
"Manifest references unconfigured source(s): " + ", ".join(missing)
|
"Manifest references unconfigured source(s): " + ", ".join(missing)
|
||||||
)
|
)
|
||||||
|
|
||||||
pending = [
|
counts = await self._jobs.counts_by_status()
|
||||||
source_id
|
pending = counts.get("queued", 0) + counts.get("claimed", 0)
|
||||||
for source_id in sorted(manifest_sources)
|
|
||||||
if await self._jobs.has_pending(source_id)
|
|
||||||
]
|
|
||||||
if pending:
|
if pending:
|
||||||
await self._pollers.close_sources()
|
await self._pollers.close_sources()
|
||||||
raise ValueError(
|
raise ValueError(
|
||||||
"Cannot replay manifest while source(s) have pending work: "
|
"Cannot replay manifest while the queue has pending work: "
|
||||||
+ ", ".join(pending)
|
f"{pending} queued/claimed job(s)"
|
||||||
)
|
)
|
||||||
|
|
||||||
seen: set[tuple[str, str]] = set()
|
seen: set[tuple[str, str]] = set()
|
||||||
|
|
|
||||||
|
|
@ -229,7 +229,7 @@ def run_batch(
|
||||||
source's sweep does not complete."""
|
source's sweep does not complete."""
|
||||||
if manifest is not None and dry_run:
|
if manifest is not None and dry_run:
|
||||||
raise typer.BadParameter("--manifest cannot be combined with --dry-run")
|
raise typer.BadParameter("--manifest cannot be combined with --dry-run")
|
||||||
if manifest is not None and output is not None:
|
if output is not None and not dry_run:
|
||||||
raise typer.BadParameter("--output is only valid with --dry-run")
|
raise typer.BadParameter("--output is only valid with --dry-run")
|
||||||
asyncio.run(
|
asyncio.run(
|
||||||
_run_batch(
|
_run_batch(
|
||||||
|
|
|
||||||
|
|
@ -237,6 +237,16 @@ def test_run_batch_manifest_conflicts_with_output(tmp_path):
|
||||||
assert "--output is only valid with --dry-run" in result.output
|
assert "--output is only valid with --dry-run" in result.output
|
||||||
|
|
||||||
|
|
||||||
|
def test_run_batch_output_requires_dry_run(tmp_path):
|
||||||
|
result = runner.invoke(
|
||||||
|
cli,
|
||||||
|
["run-batch", "--output", str(tmp_path / "out.yaml")],
|
||||||
|
)
|
||||||
|
|
||||||
|
assert result.exit_code != 0
|
||||||
|
assert "--output is only valid with --dry-run" in result.output
|
||||||
|
|
||||||
|
|
||||||
# --- serve ---
|
# --- serve ---
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -397,6 +397,36 @@ async def test_run_batch_from_manifest_rejects_pending_work(tmp_path, use_client
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_run_batch_from_manifest_rejects_unrelated_pending_work(
|
||||||
|
tmp_path, use_client
|
||||||
|
):
|
||||||
|
(tmp_path / "a.md").write_text("hello")
|
||||||
|
config = _config(tmp_path)
|
||||||
|
client = _mock_client()
|
||||||
|
use_client(client)
|
||||||
|
engine = await open_queue(config.ingester.queue)
|
||||||
|
try:
|
||||||
|
jobs = JobRepo(engine)
|
||||||
|
await jobs.enqueue("other", "file:///outside.md", op=JobOp.UPSERT)
|
||||||
|
finally:
|
||||||
|
await engine.dispose()
|
||||||
|
|
||||||
|
with pytest.raises(ValueError, match="queue has pending work"):
|
||||||
|
await IngesterApp(
|
||||||
|
config=config, db_path=tmp_path / "db.lancedb"
|
||||||
|
).run_batch_from_manifest(
|
||||||
|
_manifest(
|
||||||
|
BatchChange(
|
||||||
|
op=JobOp.UPSERT,
|
||||||
|
source_id="local",
|
||||||
|
uri=(tmp_path / "a.md").as_uri(),
|
||||||
|
discovered_at=datetime.now(UTC),
|
||||||
|
)
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_run_batch_from_manifest_rejects_duplicate_changes(tmp_path, use_client):
|
async def test_run_batch_from_manifest_rejects_duplicate_changes(tmp_path, use_client):
|
||||||
path = tmp_path / "a.md"
|
path = tmp_path / "a.md"
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue