Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
52 commits
Select commit Hold shift + click to select a range
2e64c20
feat: add persisted logs
leafty Jul 13, 2026
6b6525a
wip: forward logs from loki
leafty Jul 14, 2026
2b812b3
fix: import persisted_logs module
leafty Jul 14, 2026
8046d0e
wip: can produce unsaved log lines
leafty Jul 15, 2026
bd511ec
wip: collect session logs
leafty Jul 15, 2026
702c72f
fix config
leafty Jul 15, 2026
42c83e6
fix get_latest_log_timestamp()
leafty Jul 15, 2026
9437729
fixes
leafty Jul 15, 2026
7024719
logs
leafty Jul 15, 2026
3262663
enable loop
leafty Jul 16, 2026
47a57d7
query logs only from the last one
leafty Jul 16, 2026
92d9cbe
fix: use start = ts
leafty Jul 16, 2026
cdcaf0e
handle log bursts (?)
leafty Jul 16, 2026
6c997ca
wip: get persisted logs from the API
leafty Jul 20, 2026
e5289bb
fix get_session_logs()
leafty Jul 20, 2026
d57e393
wip get session logs
leafty Jul 20, 2026
5e94887
oops
leafty Jul 20, 2026
6dfbfb8
wip: return logs
leafty Jul 21, 2026
ab7beb0
fix ULID serialization
leafty Jul 21, 2026
a5bad1f
temp json
leafty Jul 21, 2026
dd2fdd8
wip: API
leafty Jul 21, 2026
8650392
another fix
leafty Jul 21, 2026
eb0ce2c
fix api spec
leafty Jul 21, 2026
9aa8267
fix nanotimestamp serializing
leafty Jul 21, 2026
48eb6c3
temp fix for jobs
leafty Jul 21, 2026
2d0fbce
add query params
leafty Jul 21, 2026
af2b3cb
add runs API endpoint
leafty Jul 21, 2026
83f13be
oops
leafty Jul 21, 2026
325224b
oops typo
leafty Jul 21, 2026
b893809
remove debug logs
leafty Jul 22, 2026
e43b37d
use real run_id
leafty Jul 22, 2026
4f7ba5a
Merge branch 'main' into leafty/feat-peristed-logs
leafty Jul 22, 2026
2866c94
remove debug logs
leafty Jul 24, 2026
94c3d1a
feat: purge old logs
leafty Jul 24, 2026
d82a353
fix cutoff
leafty Jul 24, 2026
62cff12
remove old runs; backtrack more
leafty Jul 24, 2026
56301a2
wip: image build logs
leafty Jul 28, 2026
01bfd98
wip: image build logs
leafty Jul 28, 2026
3cfff93
draft: persisted build logs
leafty Jul 28, 2026
6390569
improve SessionRun model
leafty Jul 30, 2026
9dedb3b
improve more models
leafty Jul 30, 2026
8a4cbb0
rename orm
leafty Jul 30, 2026
45ab475
factor db.py code
leafty Jul 30, 2026
b41fdc2
fix timestamp
leafty Jul 30, 2026
b1bb7f9
Merge branch 'main' into leafty/feat-peristed-logs
leafty Jul 30, 2026
d55a41d
implement strict check for private builds
leafty Jul 30, 2026
ba06d3c
small fixes
leafty Jul 31, 2026
22bb3f6
add test for collector
leafty Jul 31, 2026
71c4055
wip: db tests
leafty Jul 31, 2026
009cf51
wip: tests
leafty Jul 31, 2026
f6fd002
wip: db tests
leafty Aug 3, 2026
1585f41
done: db tests
leafty Aug 3, 2026
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: 2 additions & 0 deletions .vscode/settings.json
Original file line number Diff line number Diff line change
Expand Up @@ -21,4 +21,6 @@
"bases",
"components"
],
"python-envs.defaultEnvManager": "ms-python.python:poetry",
"python-envs.defaultPackageManager": "ms-python.python:poetry",
}
1 change: 1 addition & 0 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,7 @@ API_SPECS := \
components/renku_data_services/notifications/apispec.py \
components/renku_data_services/capacity_reservation/apispec.py \
components/renku_data_services/resource_usage/apispec.py \
components/renku_data_services/persisted_logs/apispec.py \
components/renku_data_services/authn/api/apispec.py

schemas: ${API_SPECS} ## Generate pydantic classes from apispec yaml files
Expand Down
11 changes: 11 additions & 0 deletions bases/renku_data_services/data_api/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
from renku_data_services.namespace.blueprints import GroupsBP
from renku_data_services.notebooks.blueprints import NotebooksNewBP
from renku_data_services.notifications.blueprints import NotificationsBP
from renku_data_services.persisted_logs.blueprints import PersistedLogsBP
from renku_data_services.platform.blueprints import PlatformConfigBP, PlatformUrlRedirectBP
from renku_data_services.project.blueprints import ProjectsBP, ProjectSessionSecretBP
from renku_data_services.repositories.blueprints import RepositoriesBP
Expand Down Expand Up @@ -301,6 +302,14 @@ def register_all_handlers(app: Sanic, dm: DependencyManager) -> Sanic:
authenticator=dm.authenticator,
rp_repo=dm.rp_repo,
)
persisted_logs = PersistedLogsBP(
name="persisted_logs",
url_prefix=url_prefix,
session_logs_repo=dm.session_logs_repo,
build_logs_repo=dm.build_logs_repo,
authenticator=dm.authenticator,
session_maker=dm.config.db.async_session_maker,
)
internal_authentication = InternalAuthenticationBP(
name="internal_authentication",
url_prefix=url_prefix,
Expand Down Expand Up @@ -343,6 +352,8 @@ def register_all_handlers(app: Sanic, dm: DependencyManager) -> Sanic:
)
if builds is not None:
app.blueprint(builds.blueprint())
if dm.config.persisted_logs.enabled:
app.blueprint(persisted_logs.blueprint())

# We need to patch sanic_ext as since version 24.12 they only send a string representation of errors
import sanic_ext.extras.validation.setup
Expand Down
8 changes: 7 additions & 1 deletion bases/renku_data_services/data_api/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
from renku_data_services.data_connectors.config import DepositConfig
from renku_data_services.db_config.config import DBConfig
from renku_data_services.notebooks.config import NotebooksConfig
from renku_data_services.persisted_logs.config import PersistedLogsConfig
from renku_data_services.secrets.config import PublicSecretsConfig
from renku_data_services.session.config import BuildsConfig
from renku_data_services.solr.solr_client import SolrClientConfig
Expand Down Expand Up @@ -48,6 +49,7 @@ class Config:
version: str
alertmanager_webhook_role: str
deposit_config: DepositConfig
persisted_logs: PersistedLogsConfig

@classmethod
def from_env(cls, db: DBConfig | None = None) -> Self:
Expand All @@ -73,11 +75,14 @@ def from_env(cls, db: DBConfig | None = None) -> Self:
gitlab_url = None

nb_config = NotebooksConfig.from_env(db, authz_config, enable_internal_gitlab=enable_internal_gitlab)

k8s_namespace = os.environ.get("K8S_NAMESPACE", "default")

return cls(
enable_internal_gitlab=enable_internal_gitlab,
version=os.environ.get("VERSION", "0.0.1"),
dummy_stores=dummy_stores,
k8s_namespace=os.environ.get("K8S_NAMESPACE", "default"),
k8s_namespace=k8s_namespace,
k8s_config_root=os.environ.get("K8S_CONFIGS_ROOT", "/secrets/kube_configs"),
db=db,
builds=BuildsConfig.from_env(),
Expand All @@ -95,4 +100,5 @@ def from_env(cls, db: DBConfig | None = None) -> Self:
log_cfg=LoggingConfig.from_env(),
alertmanager_webhook_role=os.environ.get("ALERTMANAGER_WEBHOOK_ROLE", "alertmanager-webhook"),
deposit_config=DepositConfig.from_env(nb_config.sessions.renku_url),
persisted_logs=PersistedLogsConfig.from_env(namespace=k8s_namespace),
)
15 changes: 15 additions & 0 deletions bases/renku_data_services/data_api/dependencies.py
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,10 @@
from renku_data_services.notebooks.data_sources import DataSourceRepository
from renku_data_services.notebooks.image_check import ImageCheckRepository
from renku_data_services.notifications.db import NotificationsRepository
from renku_data_services.persisted_logs.db import (
AmaltheaSessionPersistedLogsReadRepository,
ImageBuildPersistedLogsReadRepository,
)
from renku_data_services.platform.db import PlatformRepository, UrlRedirectRepository
from renku_data_services.project.db import (
ProjectMemberRepository,
Expand Down Expand Up @@ -169,6 +173,8 @@ class DependencyManager:
occurrence_repo: OccurrenceRepository
resource_requests_repo: ResourceRequestsRepo
resource_usage_service: ResourceUsageService
session_logs_repo: AmaltheaSessionPersistedLogsReadRepository
build_logs_repo: ImageBuildPersistedLogsReadRepository
zenodo_client: ZenodoAPIClient
envidat_client: EnvidatClient
job_client: DepositUploadJobClient
Expand Down Expand Up @@ -205,6 +211,7 @@ def load_apispec() -> dict[str, Any]:
renku_data_services.notifications.__file__,
renku_data_services.capacity_reservation.__file__,
renku_data_services.resource_usage.__file__,
renku_data_services.persisted_logs.__file__,
renku_data_services.authn.api.__file__,
]

Expand Down Expand Up @@ -462,6 +469,12 @@ def from_env(cls) -> DependencyManager:
occurrence_repo = OccurrenceRepository(
session_maker=config.db.async_session_maker,
)
session_logs_repo = AmaltheaSessionPersistedLogsReadRepository(authz=authz)
build_logs_repo = ImageBuildPersistedLogsReadRepository(
authz=authz,
builds_config=config.builds,
git_repositories_repo=git_repositories_repo,
)
return cls(
config,
k8s_client=client,
Expand Down Expand Up @@ -507,6 +520,8 @@ def from_env(cls) -> DependencyManager:
occurrence_repo=occurrence_repo,
resource_requests_repo=resource_requests_repo,
resource_usage_service=resource_usage_service,
session_logs_repo=session_logs_repo,
build_logs_repo=build_logs_repo,
zenodo_client=ZenodoAPIClient(),
envidat_client=EnvidatClient(),
job_client=job_client,
Expand Down
4 changes: 4 additions & 0 deletions bases/renku_data_services/data_tasks/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
from renku_data_services.app_config.config import KeycloakConfig
from renku_data_services.authz.config import AuthzConfig
from renku_data_services.db_config.config import DBConfig
from renku_data_services.persisted_logs.config import PersistedLogsConfig
from renku_data_services.solr.solr_client import SolrClientConfig


Expand Down Expand Up @@ -42,6 +43,7 @@ class Config:
posthog: PosthogConfig
authz: AuthzConfig
keycloak: KeycloakConfig | None
persisted_logs: PersistedLogsConfig
k8s_config_root: str
dummy_stores: bool
max_retry_wait_seconds: int
Expand Down Expand Up @@ -77,6 +79,7 @@ def from_env(cls) -> Config:
session_quota_alert_remaining_threshold = int(os.environ.get("SESSION_QUOTA_ALERT_REMAINING_THRESHOLD_P", 20))
session_quota_alert_critical = int(os.environ.get("SESSION_QUOTA_ALERT_CRITICAL_M", 10))

k8s_namespace = os.environ.get("KUBERNETES_NAMESPACE", os.environ.get("K8S_NAMESPACE", "default"))
k8s_config_root = os.environ.get("K8S_CONFIG_ROOT", "/secrets/kube_configs")

enable_resource_request_tracking = os.environ.get("ENABLE_RESOURCE_REQUEST_TRACKING", "false").lower() == "true"
Expand All @@ -91,6 +94,7 @@ def from_env(cls) -> Config:
posthog=posthog_config,
authz=authz,
keycloak=keycloak,
persisted_logs=PersistedLogsConfig.from_env(namespace=k8s_namespace),
k8s_config_root=k8s_config_root,
tcp_host=tcp_host,
tcp_port=tcp_port,
Expand Down
8 changes: 8 additions & 0 deletions bases/renku_data_services/data_tasks/dependencies.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
from renku_data_services.namespace.db import GroupRepository
from renku_data_services.notebooks.constants import AMALTHEA_SESSION_GVK
from renku_data_services.notifications.db import NotificationsRepository
from renku_data_services.persisted_logs.collector import PersistedLogsCollector
from renku_data_services.project.db import ProjectRepository
from renku_data_services.resource_usage.core import (
DefaultResourcesRequestRecorder,
Expand Down Expand Up @@ -56,6 +57,7 @@ class DependencyManager:
notifications_repo: NotificationsRepository
resource_usage_service: ResourceUsageService
resource_requests_repo: ResourceRequestsRepo
persisted_logs_collector: PersistedLogsCollector

@classmethod
def from_env(cls, cfg: Config | None = None) -> "DependencyManager":
Expand Down Expand Up @@ -151,6 +153,11 @@ def from_env(cls, cfg: Config | None = None) -> "DependencyManager":
realm=cfg.keycloak.realm,
)

persisted_logs_collector = PersistedLogsCollector.from_config(
config=cfg.persisted_logs,
session_maker=cfg.db.async_session_maker,
)

return cls(
config=cfg,
search_updates_repo=search_updates_repo,
Expand All @@ -167,4 +174,5 @@ def from_env(cls, cfg: Config | None = None) -> "DependencyManager":
notifications_repo=notifications_repo,
resource_usage_service=resource_usage_service,
resource_requests_repo=resource_requests_repo,
persisted_logs_collector=persisted_logs_collector,
)
24 changes: 24 additions & 0 deletions bases/renku_data_services/data_tasks/task_defs.py
Original file line number Diff line number Diff line change
Expand Up @@ -577,6 +577,28 @@ async def monitor_session_quota_and_send_alerts(dm: DependencyManager) -> None:
await asyncio.sleep(dm.config.session_quota_alert_check_interval_s)


async def collect_persisted_logs(dm: DependencyManager) -> None:
"""Collect persisted logs from Loki."""
while True:
try:
await dm.persisted_logs_collector.collect_persisted_logs()
except Exception as e:
logger.warning(f"Failed to collect persisted logs: {e}", exc_info=True)
else:
await asyncio.sleep(1)


async def purge_expired_persisted_logs(dm: DependencyManager) -> None:
"""Purge expired persisted logs from the database."""
while True:
try:
await dm.persisted_logs_collector.purge_expired_logs()
except Exception as e:
logger.warning(f"Failed to purge expired persisted logs: {e}", exc_info=True)
else:
await asyncio.sleep(dm.config.long_task_period_s)


def all_tasks(dm: DependencyManager) -> TaskDefininions:
"""A dict of task factories to be managed in main."""
# Impl. note: We pass the entire config to the coroutines, because
Expand All @@ -603,5 +625,7 @@ def all_tasks(dm: DependencyManager) -> TaskDefininions:
"cleanup_orphaned_capacity_reservations": lambda: cleanup_orphaned_capacity_reservations(dm),
"record_resource_requests": lambda: record_resource_requests(dm),
"monitor_session_quota_and_send_alerts": lambda: monitor_session_quota_and_send_alerts(dm),
"collect_persisted_logs": lambda: collect_persisted_logs(dm),
"purge_expired_persisted_logs": lambda: purge_expired_persisted_logs(dm),
}
)
2 changes: 2 additions & 0 deletions components/renku_data_services/migrations/env.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
from renku_data_services.migrations.utils import run_migrations
from renku_data_services.namespace.orm import BaseORM as namespaces
from renku_data_services.notifications.orm import BaseORM as notifications
from renku_data_services.persisted_logs.orm import BaseORM as persisted_logs
from renku_data_services.platform.orm import BaseORM as platform
from renku_data_services.project.orm import BaseORM as project
from renku_data_services.resource_usage.orm import BaseORM as resource_usage
Expand All @@ -30,6 +31,7 @@
metrics.metadata,
namespaces.metadata,
notifications.metadata,
persisted_logs.metadata,
platform.metadata,
project.metadata,
search.metadata,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
"""wip: redefine tables

Revision ID: 01180f797019
Revises: 2537a8e1df45
Create Date: 2026-07-15 12:50:59.650960

"""

import sqlalchemy as sa
from alembic import op

from renku_data_services.utils.sqlalchemy import ULIDType

# revision identifiers, used by Alembic.
revision = "01180f797019"
down_revision = "2537a8e1df45"
branch_labels = None
depends_on = None


def upgrade() -> None:
# ### commands auto generated by Alembic - please adjust! ###
op.create_table(
"session_runs",
sa.Column("id", ULIDType(), nullable=False),
sa.Column("user_id", sa.String(length=36), nullable=False),
sa.Column("launch_id", sa.String(), nullable=False),
sa.Column("launcher_id", ULIDType(), nullable=False),
sa.Column("submission_id", sa.String(), nullable=True),
sa.ForeignKeyConstraint(
["launcher_id"],
["sessions.launchers.id"],
),
sa.ForeignKeyConstraint(
["user_id"],
["users.users.keycloak_id"],
),
sa.PrimaryKeyConstraint("id"),
schema="persisted_logs",
)
op.create_index(
op.f("ix_persisted_logs_session_runs_launcher_id"),
"session_runs",
["launcher_id"],
unique=False,
schema="persisted_logs",
)
op.create_index(
op.f("ix_persisted_logs_session_runs_user_id"),
"session_runs",
["user_id"],
unique=False,
schema="persisted_logs",
)
op.create_table(
"amalthea_session_logs",
sa.Column("id", sa.String(), nullable=False),
sa.Column("run_id", ULIDType(), nullable=False),
sa.Column("container", sa.String(), nullable=False),
sa.Column("timestamp", sa.BigInteger(), nullable=False),
sa.Column("log_line", sa.String(), nullable=False),
sa.ForeignKeyConstraint(["run_id"], ["persisted_logs.session_runs.id"], ondelete="CASCADE"),
sa.PrimaryKeyConstraint("id"),
schema="persisted_logs",
)
op.create_index(
op.f("ix_persisted_logs_amalthea_session_logs_run_id"),
"amalthea_session_logs",
["run_id"],
unique=False,
schema="persisted_logs",
)
op.create_index(
op.f("ix_persisted_logs_amalthea_session_logs_timestamp"),
"amalthea_session_logs",
["timestamp"],
unique=False,
schema="persisted_logs",
)
# ### end Alembic commands ###


def downgrade() -> None:
# ### commands auto generated by Alembic - please adjust! ###
op.drop_index(
op.f("ix_persisted_logs_amalthea_session_logs_timestamp"),
table_name="amalthea_session_logs",
schema="persisted_logs",
)
op.drop_index(
op.f("ix_persisted_logs_amalthea_session_logs_run_id"),
table_name="amalthea_session_logs",
schema="persisted_logs",
)
op.drop_table("amalthea_session_logs", schema="persisted_logs")
op.drop_index(op.f("ix_persisted_logs_session_runs_user_id"), table_name="session_runs", schema="persisted_logs")
op.drop_index(
op.f("ix_persisted_logs_session_runs_launcher_id"), table_name="session_runs", schema="persisted_logs"
)
op.drop_table("session_runs", schema="persisted_logs")
# ### end Alembic commands ###
Loading
Loading