Skip to content
Draft
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
74 changes: 71 additions & 3 deletions poetry.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,7 @@ pytest-retry = "1.7.0"
datamodel-code-generator = {extras = ["http"], version = ">=0.54,<0.65"}
pytest-asyncio = "^1.3.0"
pytest-cov = "^7.1.0"
pylint = "^4.0.8"

[tool.poetry.group.docs.dependencies]
sphinx = "^5.1.1"
Expand Down
20 changes: 14 additions & 6 deletions pysus/management/catalog.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
from __future__ import annotations

import hashlib
from dataclasses import dataclass
from datetime import datetime
from pathlib import Path
from typing import TYPE_CHECKING, Any
Expand Down Expand Up @@ -64,6 +65,15 @@ def sha256_of(path: Path) -> str:
return digest.hexdigest()


@dataclass(frozen=True)
class FileOriginMeta:
"""Encapsulates file origin metadata for updates."""

modified: datetime | None
size: int
source_sha256: str | None = None


class CatalogWriter:
"""Upsert dataset/group/file/column metadata into the DuckLake catalogs."""

Expand Down Expand Up @@ -228,22 +238,20 @@ def touch_file(
self,
cursor,
file_id: int,
origin_modified: datetime | None,
origin_size: int,
source_sha256: str | None = None,
meta: FileOriginMeta,
) -> None:
"""Update origin metadata without replacing the artifact."""
if source_sha256 is not None:
if meta.source_sha256 is not None:
cursor.execute(
"UPDATE pysus.files SET origin_modified = ?, "
"origin_size = ?, source_sha256 = ? WHERE id = ?",
(origin_modified, origin_size, source_sha256, file_id),
(meta.modified, meta.size, meta.source_sha256, file_id),
)
else:
cursor.execute(
"UPDATE pysus.files SET origin_modified = ?, "
"origin_size = ? WHERE id = ?",
(origin_modified, origin_size, file_id),
(meta.modified, meta.size, file_id),
)

def delete_file(self, cursor, file_id: int) -> None:
Expand Down
16 changes: 10 additions & 6 deletions pysus/management/sync.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@
from pysus.api.errors import AuthenticationError, ConnectionError
from pysus.api.models import BaseRemoteFile

from .catalog import CatalogWriter, sha256_of
from .catalog import CatalogWriter, FileOriginMeta, sha256_of
from .compare import Comparator
from .inventory import Inventory
from .records import (
Expand Down Expand Up @@ -365,8 +365,10 @@ async def upload_file(
writer.touch_file(
dataset_cursor,
file_id,
self._safe_modify(file),
self._safe_size(file),
FileOriginMeta(
modified=self._safe_modify(file),
size=self._safe_size(file),
)
)
dataset_conn.commit()
dataset_cursor.execute("CHECKPOINT")
Expand Down Expand Up @@ -394,9 +396,11 @@ async def upload_file(
writer.touch_file(
dataset_cursor,
file_id,
self._safe_modify(file),
self._safe_size(file),
source_sha256=raw_digest,
FileOriginMeta(
modified=self._safe_modify(file),
size=self._safe_size(file),
source_sha256=raw_digest,
)
)
dataset_conn.commit()
dataset_cursor.execute("CHECKPOINT")
Expand Down
19 changes: 16 additions & 3 deletions pysus/tests/management/test_catalog.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
import duckdb
import pyarrow as pa
import pytest
from pysus.management.catalog import CatalogWriter
from pysus.management.catalog import CatalogWriter, FileOriginMeta

_SCHEMA = """
CREATE SCHEMA pysus;
Expand Down Expand Up @@ -120,7 +120,13 @@ def test_updates_origin_metadata(self, writer_and_cursor):
writer, cursor, _ = writer_and_cursor
_insert(cursor, "public/data/x.parquet")
writer.touch_file(
cursor, 1, datetime(2026, 2, 2), 99, source_sha256="cc" * 32
cursor,
1,
FileOriginMeta(
modified=datetime(2026, 2, 2),
size=99,
source_sha256="cc" * 32,
)
)
result = writer.get_file_full(cursor, "public/data/x.parquet")
assert result is not None
Expand All @@ -132,7 +138,14 @@ def test_updates_origin_metadata(self, writer_and_cursor):
def test_without_source_sha256(self, writer_and_cursor):
writer, cursor, _ = writer_and_cursor
_insert(cursor, "public/data/x.parquet")
writer.touch_file(cursor, 1, datetime(2026, 2, 2), 99)
writer.touch_file(
cursor,
1,
FileOriginMeta(
modified=datetime(2026, 2, 2),
size=99,
)
)
result = writer.get_file_full(cursor, "public/data/x.parquet")
assert result is not None
assert result[1] == datetime(2026, 2, 2)
Expand Down