diff --git a/.vscode/settings.json b/.vscode/settings.json index 32c4a1c129..b5d233914a 100644 --- a/.vscode/settings.json +++ b/.vscode/settings.json @@ -21,4 +21,6 @@ "bases", "components" ], + "python-envs.defaultEnvManager": "ms-python.python:poetry", + "python-envs.defaultPackageManager": "ms-python.python:poetry", } diff --git a/Makefile b/Makefile index d50c138419..3e607162f9 100644 --- a/Makefile +++ b/Makefile @@ -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 diff --git a/bases/renku_data_services/data_api/app.py b/bases/renku_data_services/data_api/app.py index 5d3a23cf16..781f40786b 100644 --- a/bases/renku_data_services/data_api/app.py +++ b/bases/renku_data_services/data_api/app.py @@ -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 @@ -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, @@ -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 diff --git a/bases/renku_data_services/data_api/config.py b/bases/renku_data_services/data_api/config.py index 045176de5b..c1ec557b23 100644 --- a/bases/renku_data_services/data_api/config.py +++ b/bases/renku_data_services/data_api/config.py @@ -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 @@ -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: @@ -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(), @@ -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), ) diff --git a/bases/renku_data_services/data_api/dependencies.py b/bases/renku_data_services/data_api/dependencies.py index 6c4200a39f..4c753bdcdd 100644 --- a/bases/renku_data_services/data_api/dependencies.py +++ b/bases/renku_data_services/data_api/dependencies.py @@ -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, @@ -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 @@ -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__, ] @@ -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, @@ -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, diff --git a/bases/renku_data_services/data_tasks/config.py b/bases/renku_data_services/data_tasks/config.py index 7d09340652..7c0f4d03ac 100644 --- a/bases/renku_data_services/data_tasks/config.py +++ b/bases/renku_data_services/data_tasks/config.py @@ -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 @@ -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 @@ -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" @@ -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, diff --git a/bases/renku_data_services/data_tasks/dependencies.py b/bases/renku_data_services/data_tasks/dependencies.py index c7a9ac0360..26d9f42f67 100644 --- a/bases/renku_data_services/data_tasks/dependencies.py +++ b/bases/renku_data_services/data_tasks/dependencies.py @@ -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, @@ -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": @@ -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, @@ -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, ) diff --git a/bases/renku_data_services/data_tasks/task_defs.py b/bases/renku_data_services/data_tasks/task_defs.py index 903c2a9773..e762fe8519 100644 --- a/bases/renku_data_services/data_tasks/task_defs.py +++ b/bases/renku_data_services/data_tasks/task_defs.py @@ -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 @@ -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), } ) diff --git a/components/renku_data_services/migrations/env.py b/components/renku_data_services/migrations/env.py index ebe4f15665..812425be6f 100644 --- a/components/renku_data_services/migrations/env.py +++ b/components/renku_data_services/migrations/env.py @@ -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 @@ -30,6 +31,7 @@ metrics.metadata, namespaces.metadata, notifications.metadata, + persisted_logs.metadata, platform.metadata, project.metadata, search.metadata, diff --git a/components/renku_data_services/migrations/versions/01180f797019_wip_redefine_tables.py b/components/renku_data_services/migrations/versions/01180f797019_wip_redefine_tables.py new file mode 100644 index 0000000000..f884dac00f --- /dev/null +++ b/components/renku_data_services/migrations/versions/01180f797019_wip_redefine_tables.py @@ -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 ### diff --git a/components/renku_data_services/migrations/versions/2537a8e1df45_wip.py b/components/renku_data_services/migrations/versions/2537a8e1df45_wip.py new file mode 100644 index 0000000000..5ee54f5c96 --- /dev/null +++ b/components/renku_data_services/migrations/versions/2537a8e1df45_wip.py @@ -0,0 +1,74 @@ +"""wip + +Revision ID: 2537a8e1df45 +Revises: eadfb5e7e7cb +Create Date: 2026-07-15 11:42:06.232593 + +""" + +import sqlalchemy as sa +from alembic import op +from sqlalchemy.dialects import postgresql + +# revision identifiers, used by Alembic. +revision = "2537a8e1df45" +down_revision = "eadfb5e7e7cb" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.drop_index( + "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("ix_persisted_logs_session_runs_launcher_id", table_name="session_runs", schema="persisted_logs") + op.drop_index("ix_persisted_logs_session_runs_user_id", table_name="session_runs", schema="persisted_logs") + op.drop_table("session_runs", schema="persisted_logs") + + +def downgrade() -> None: + op.create_table( + "session_runs", + sa.Column("id", sa.VARCHAR(), autoincrement=False, nullable=False), + sa.Column("user_id", sa.VARCHAR(length=36), autoincrement=False, nullable=False), + sa.Column("launch_id", sa.VARCHAR(), autoincrement=False, nullable=False), + sa.Column("launcher_id", sa.VARCHAR(), autoincrement=False, nullable=False), + sa.Column("submission_id", sa.VARCHAR(), autoincrement=False, nullable=True), + sa.Column("first_log", postgresql.TIMESTAMP(timezone=True), autoincrement=False, nullable=False), + sa.Column("last_log", postgresql.TIMESTAMP(timezone=True), autoincrement=False, nullable=False), + sa.ForeignKeyConstraint(["launcher_id"], ["sessions.launchers.id"], name="session_runs_launcher_id_fkey"), + sa.ForeignKeyConstraint(["user_id"], ["users.users.keycloak_id"], name="session_runs_user_id_fkey"), + sa.PrimaryKeyConstraint("id", name="session_runs_pkey"), + schema="persisted_logs", + ) + op.create_index( + "ix_persisted_logs_session_runs_user_id", "session_runs", ["user_id"], unique=False, schema="persisted_logs" + ) + op.create_index( + "ix_persisted_logs_session_runs_launcher_id", + "session_runs", + ["launcher_id"], + unique=False, + schema="persisted_logs", + ) + op.create_table( + "amalthea_session_logs", + sa.Column("id", sa.VARCHAR(), server_default=sa.text("generate_ulid()"), autoincrement=False, nullable=False), + sa.Column("run_id", sa.VARCHAR(), autoincrement=False, nullable=False), + sa.Column("container", sa.VARCHAR(), autoincrement=False, nullable=False), + sa.Column("timestamp", postgresql.TIMESTAMP(timezone=True), autoincrement=False, nullable=False), + sa.Column("log_line", sa.VARCHAR(), autoincrement=False, nullable=False), + sa.ForeignKeyConstraint( + ["run_id"], ["persisted_logs.session_runs.id"], name="amalthea_session_logs_run_id_fkey", ondelete="CASCADE" + ), + sa.PrimaryKeyConstraint("id", name="amalthea_session_logs_pkey"), + schema="persisted_logs", + ) + op.create_index( + "ix_persisted_logs_amalthea_session_logs_run_id", + "amalthea_session_logs", + ["run_id"], + unique=False, + schema="persisted_logs", + ) diff --git a/components/renku_data_services/migrations/versions/906ec89ea06f_wip_update_persisted_logs.py b/components/renku_data_services/migrations/versions/906ec89ea06f_wip_update_persisted_logs.py new file mode 100644 index 0000000000..2908836e35 --- /dev/null +++ b/components/renku_data_services/migrations/versions/906ec89ea06f_wip_update_persisted_logs.py @@ -0,0 +1,31 @@ +"""wip: update persisted logs + +Revision ID: 906ec89ea06f +Revises: 01180f797019 +Create Date: 2026-07-22 08:50:21.380775 + +""" + +import sqlalchemy as sa +from alembic import op + +# revision identifiers, used by Alembic. +revision = "906ec89ea06f" +down_revision = "01180f797019" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.execute("DELETE FROM persisted_logs.session_runs") + op.add_column("session_runs", sa.Column("session_uid", sa.String(), nullable=True), schema="persisted_logs") + op.drop_column("session_runs", "launch_id", schema="persisted_logs") + + +def downgrade() -> None: + op.add_column( + "session_runs", + sa.Column("launch_id", sa.VARCHAR(), autoincrement=False, nullable=False), + schema="persisted_logs", + ) + op.drop_column("session_runs", "session_uid", schema="persisted_logs") diff --git a/components/renku_data_services/migrations/versions/eadfb5e7e7cb_feat_add_persisted_logs.py b/components/renku_data_services/migrations/versions/eadfb5e7e7cb_feat_add_persisted_logs.py new file mode 100644 index 0000000000..6103710e33 --- /dev/null +++ b/components/renku_data_services/migrations/versions/eadfb5e7e7cb_feat_add_persisted_logs.py @@ -0,0 +1,91 @@ +"""feat: add persisted logs + +Revision ID: eadfb5e7e7cb +Revises: 0a6cee40fe0d +Create Date: 2026-07-13 12:35:55.042555 + +""" + +import sqlalchemy as sa +from alembic import op + +from renku_data_services.utils.sqlalchemy import ULIDType + +# revision identifiers, used by Alembic. +revision = "eadfb5e7e7cb" +down_revision = "0a6cee40fe0d" +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.Column("first_log", sa.DateTime(timezone=True), nullable=False), + sa.Column("last_log", sa.DateTime(timezone=True), nullable=False), + 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", ULIDType(), server_default=sa.text("generate_ulid()"), nullable=False), + sa.Column("run_id", ULIDType(), nullable=False), + sa.Column("container", sa.String(), nullable=False), + sa.Column("timestamp", sa.DateTime(timezone=True), 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", + ) + # ### end Alembic commands ### + + +def downgrade() -> None: + # ### commands auto generated by Alembic - please adjust! ### + 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 ### diff --git a/components/renku_data_services/migrations/versions/f65223ca0978_feat_add_persisted_build_logs.py b/components/renku_data_services/migrations/versions/f65223ca0978_feat_add_persisted_build_logs.py new file mode 100644 index 0000000000..497aec5a8c --- /dev/null +++ b/components/renku_data_services/migrations/versions/f65223ca0978_feat_add_persisted_build_logs.py @@ -0,0 +1,56 @@ +"""feat: add persisted build logs + +Revision ID: f65223ca0978 +Revises: 906ec89ea06f +Create Date: 2026-07-28 07:59:44.706230 + +""" + +import sqlalchemy as sa +from alembic import op + +from renku_data_services.utils.sqlalchemy import ULIDType + +# revision identifiers, used by Alembic. +revision = "f65223ca0978" +down_revision = "906ec89ea06f" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.create_table( + "image_build_logs", + sa.Column("id", sa.String(), nullable=False), + sa.Column("build_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(["build_id"], ["sessions.builds.id"], ondelete="CASCADE"), + sa.PrimaryKeyConstraint("id"), + schema="persisted_logs", + ) + op.create_index( + op.f("ix_persisted_logs_image_build_logs_build_id"), + "image_build_logs", + ["build_id"], + unique=False, + schema="persisted_logs", + ) + op.create_index( + op.f("ix_persisted_logs_image_build_logs_timestamp"), + "image_build_logs", + ["timestamp"], + unique=False, + schema="persisted_logs", + ) + + +def downgrade() -> None: + op.drop_index( + op.f("ix_persisted_logs_image_build_logs_timestamp"), table_name="image_build_logs", schema="persisted_logs" + ) + op.drop_index( + op.f("ix_persisted_logs_image_build_logs_build_id"), table_name="image_build_logs", schema="persisted_logs" + ) + op.drop_table("image_build_logs", schema="persisted_logs") diff --git a/components/renku_data_services/persisted_logs/__init__.py b/components/renku_data_services/persisted_logs/__init__.py new file mode 100644 index 0000000000..83f61e78c4 --- /dev/null +++ b/components/renku_data_services/persisted_logs/__init__.py @@ -0,0 +1,4 @@ +"""Persisted logs module. + +Provides persisted logs for user workloads: interactive sessions, offline jobs, image builds, etc. +""" diff --git a/components/renku_data_services/persisted_logs/api.spec.yaml b/components/renku_data_services/persisted_logs/api.spec.yaml new file mode 100644 index 0000000000..c9a83afbc2 --- /dev/null +++ b/components/renku_data_services/persisted_logs/api.spec.yaml @@ -0,0 +1,237 @@ +openapi: 3.0.2 +info: + title: Renku Data Services API + description: | + This service is the main backend for Renku. It provides information about users, projects, + cloud storage, access to compute resources and many other things. + version: v1 +servers: + - url: /api/data +paths: + /persisted_logs/sessions/{launcher_id}: + get: + summary: Get persisted logs for a session + description: | + Returns logs for a given session launcher belonging to the current user. + + * If the `run_id` is not specified, logs from the most recent run will be returned. + * For offline jobs (`launcher_type: non-interactive`), the `submission_id` can + also be used to get the corresponding logs. + + NOTE: logs are persisted for a limited amount of time. Once expired, logs are + purged from the database. + parameters: + - in: path + name: launcher_id + required: true + schema: + $ref: "#/components/schemas/Ulid" + - in: query + name: params + style: form + explode: true + schema: + $ref: "#/components/schemas/PersistedSessionLogsGetQuery" + responses: + "200": + description: | + The persisted logs from the corresponding session run. + content: + application/json: + schema: + $ref: "#/components/schemas/PersistedSessionLogs" + "404": + description: | + The session launcher does not exist or there are no logs persisted for it at the moment. + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + default: + $ref: "#/components/responses/Error" + tags: + - persisted_logs + /persisted_logs/sessions/{launcher_id}/runs: + get: + summary: Get the list of session runs for a given session launcher + description: | + Returns the list of session runs for the given session launcher. + + NOTE: session runs are removed when the corresponding logs have expired. + parameters: + - in: path + name: launcher_id + required: true + schema: + $ref: "#/components/schemas/Ulid" + responses: + "200": + description: | + The session runs for which logs exist. + content: + application/json: + schema: + $ref: "#/components/schemas/SessionRuns" + "404": + description: The session launcher does not exist + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + default: + $ref: "#/components/responses/Error" + tags: + - persisted_logs + /persisted_logs/builds/{build_id}: + get: + summary: Get persisted logs for a build + description: | + Returns logs for a given image build. + + NOTE: logs are persisted for a limited amount of time. Once expired, logs are + purged from the database. + parameters: + - in: path + name: build_id + required: true + schema: + $ref: "#/components/schemas/Ulid" + responses: + "200": + description: | + The image build logs from the corresponding image build. + content: + application/json: + schema: + $ref: "#/components/schemas/PersistedBuildLogs" + "404": + description: The image build does not exist + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" + default: + $ref: "#/components/responses/Error" + tags: + - persisted_logs +components: + schemas: + PersistedSessionLogs: + description: Persisted logs for a session + type: object + properties: + run: + $ref: "#/components/schemas/SessionRun" + logs: + $ref: "#/components/schemas/LogsPerContainer" + required: + - run + - logs + SessionRuns: + description: A list of session runs + type: array + items: + $ref: "#/components/schemas/SessionRun" + SessionRun: + description: A session run + type: object + properties: + id: + $ref: "#/components/schemas/Ulid" + session_uid: + type: string + launcher_id: + $ref: "#/components/schemas/Ulid" + submission_id: + type: string + required: + - id + - launcher_id + PersistedBuildLogs: + description: Persisted logs for an image build + type: object + properties: + logs: + $ref: "#/components/schemas/LogsPerContainer" + required: + - logs + LogsPerContainer: + description: Logs organized by pod container + type: array + items: + $ref: "#/components/schemas/ContainerLogs" + ContainerLogs: + description: Logs of a single container + type: object + properties: + container: + type: string + logs: + $ref: "#/components/schemas/PersistedLogLines" + required: + - container + - logs + PersistedLogLines: + description: Stream of log lines + type: array + items: + $ref: "#/components/schemas/PersistedLogLine" + PersistedLogLine: + description: A timestamped log line + type: object + properties: + timestamp: + $ref: "#/components/schemas/NanoTimestamp" + log_line: + type: string + required: + - timestamp + - log_line + Ulid: + description: ULID identifier + type: string + minLength: 26 + maxLength: 26 + pattern: "^[0-7][0-9A-HJKMNP-TV-Z]{25}$" # This is case-insensitive + NanoTimestamp: + description: UNIX timestamp in nanoseconds + type: string + pattern: "[0-9]+" + PersistedSessionLogsGetQuery: + description: Query params for querying persisted logs of a session + type: object + properties: + run_id: + $ref: "#/components/schemas/Ulid" + submission_id: + type: string + ErrorResponse: + type: object + properties: + error: + type: object + properties: + code: + type: integer + minimum: 0 + exclusiveMinimum: true + example: 1404 + detail: + type: string + example: "A more detailed optional message showing what the problem was" + message: + type: string + example: "Something went wrong - please try again later" + trace_id: + type: string + example: "ac93950e9e114a55c67fb8e5ef519bbe" + description: Sentry trace ID for linking to corresponding log entries + required: ["code", "message"] + required: ["error"] + responses: + Error: + description: The schema for all 4xx and 5xx responses + content: + "application/json": + schema: + $ref: "#/components/schemas/ErrorResponse" diff --git a/components/renku_data_services/persisted_logs/apispec.py b/components/renku_data_services/persisted_logs/apispec.py new file mode 100644 index 0000000000..84aff8f2b9 --- /dev/null +++ b/components/renku_data_services/persisted_logs/apispec.py @@ -0,0 +1,90 @@ +# generated by datamodel-codegen: +# filename: api.spec.yaml +# timestamp: 2026-07-30T11:38:47+00:00 + +from __future__ import annotations + +from pydantic import Field, RootModel +from renku_data_services.persisted_logs.apispec_base import BaseAPISpec + + +class PersistedSessionLogsGetQuery(BaseAPISpec): + run_id: str | None = Field( + None, + description="ULID identifier", + max_length=26, + min_length=26, + pattern="^[0-7][0-9A-HJKMNP-TV-Z]{25}$", + ) + submission_id: str | None = None + + +class Error(BaseAPISpec): + code: int = Field(..., examples=[1404], gt=0) + detail: str | None = Field( + None, examples=["A more detailed optional message showing what the problem was"] + ) + message: str = Field( + ..., examples=["Something went wrong - please try again later"] + ) + trace_id: str | None = Field( + None, + description="Sentry trace ID for linking to corresponding log entries", + examples=["ac93950e9e114a55c67fb8e5ef519bbe"], + ) + + +class ErrorResponse(BaseAPISpec): + error: Error + + +class PersistedLogsSessionsLauncherIdGetParametersQuery(BaseAPISpec): + params: PersistedSessionLogsGetQuery | None = None + + +class SessionRun(BaseAPISpec): + id: str = Field( + ..., + description="ULID identifier", + max_length=26, + min_length=26, + pattern="^[0-7][0-9A-HJKMNP-TV-Z]{25}$", + ) + session_uid: str | None = None + launcher_id: str = Field( + ..., + description="ULID identifier", + max_length=26, + min_length=26, + pattern="^[0-7][0-9A-HJKMNP-TV-Z]{25}$", + ) + submission_id: str | None = None + + +class PersistedLogLine(BaseAPISpec): + timestamp: str = Field( + ..., description="UNIX timestamp in nanoseconds", pattern="[0-9]+" + ) + log_line: str + + +class SessionRuns(RootModel[list[SessionRun]]): + root: list[SessionRun] = Field(..., description="A list of session runs") + + +class ContainerLogs(BaseAPISpec): + container: str + logs: list[PersistedLogLine] = Field(..., description="Stream of log lines") + + +class PersistedSessionLogs(BaseAPISpec): + run: SessionRun + logs: list[ContainerLogs] = Field( + ..., description="Logs organized by pod container" + ) + + +class PersistedBuildLogs(BaseAPISpec): + logs: list[ContainerLogs] = Field( + ..., description="Logs organized by pod container" + ) diff --git a/components/renku_data_services/persisted_logs/apispec_base.py b/components/renku_data_services/persisted_logs/apispec_base.py new file mode 100644 index 0000000000..edb200ef95 --- /dev/null +++ b/components/renku_data_services/persisted_logs/apispec_base.py @@ -0,0 +1,31 @@ +"""Base models for API specifications.""" + +from typing import Any + +from pydantic import BaseModel, ConfigDict, field_validator +from ulid import ULID + + +class BaseAPISpec(BaseModel): + """Base API specification.""" + + model_config = ConfigDict( + # Enables orm mode for pydantic.""" + from_attributes=True, + ) + + @field_validator("*", mode="before", check_fields=False) + @classmethod + def serialize_ulid(cls, value: Any) -> Any: + """Handle ULIDs.""" + if isinstance(value, ULID): + return str(value) + return value + + @field_validator("timestamp", mode="before", check_fields=False) + @classmethod + def serialize_nano_timestamp(cls, value: Any) -> Any: + """Handle serializing nanosecond timestamps to string.""" + if isinstance(value, int): + return str(value) + return value diff --git a/components/renku_data_services/persisted_logs/blueprints.py b/components/renku_data_services/persisted_logs/blueprints.py new file mode 100644 index 0000000000..bcf6f19901 --- /dev/null +++ b/components/renku_data_services/persisted_logs/blueprints.py @@ -0,0 +1,86 @@ +"""Persisted logs blueprints.""" + +from collections.abc import Callable +from dataclasses import dataclass + +from sanic import Request +from sanic.response import JSONResponse +from sqlalchemy.ext.asyncio import AsyncSession +from ulid import ULID + +from renku_data_services import base_models, errors +from renku_data_services.base_api.auth import authenticate, only_authenticated +from renku_data_services.base_api.blueprint import BlueprintFactoryResponse, CustomBlueprint +from renku_data_services.base_api.misc import validate_query +from renku_data_services.base_models.validation import validated_json +from renku_data_services.persisted_logs import apispec, models +from renku_data_services.persisted_logs.db import ( + AmaltheaSessionPersistedLogsReadRepository, + ImageBuildPersistedLogsReadRepository, +) + + +@dataclass(kw_only=True) +class PersistedLogsBP(CustomBlueprint): + """Handlers for querying persisted logs.""" + + session_logs_repo: AmaltheaSessionPersistedLogsReadRepository + build_logs_repo: ImageBuildPersistedLogsReadRepository + authenticator: base_models.Authenticator + session_maker: Callable[..., AsyncSession] + + def get_session_logs(self) -> BlueprintFactoryResponse: + """Get persisted sessions logs.""" + + @authenticate(self.authenticator) + @only_authenticated + @validate_query(query=apispec.PersistedSessionLogsGetQuery) + async def _get_session_logs( + _: Request, user: base_models.APIUser, launcher_id: ULID, query: apispec.PersistedSessionLogsGetQuery + ) -> JSONResponse: + run_id = ULID.from_str(query.run_id) if query.run_id else None + async with self.session_maker() as session, session.begin(): + result = await self.session_logs_repo.get_session_logs( + session=session, + user=user, + launcher_id=launcher_id, + run_id=run_id, + submission_id=query.submission_id, + ) + if result is None: + # TODO: adjust error message when params are passed in the query + raise errors.MissingResourceError( + message=f"Session launcher with id '{launcher_id}' does not have persisted logs." + ) + return validated_json(apispec.PersistedSessionLogs, result) + + return "/persisted_logs/sessions/", ["GET"], _get_session_logs + + def get_session_runs(self) -> BlueprintFactoryResponse: + """Get the session runs for a given session launcher.""" + + @authenticate(self.authenticator) + @only_authenticated + async def _get_session_runs(_: Request, user: base_models.APIUser, launcher_id: ULID) -> JSONResponse: + async with self.session_maker() as session, session.begin(): + session_runs = self.session_logs_repo.get_session_runs( + session=session, user=user, launcher_id=launcher_id + ) + result: list[models.SessionRun] = [] + async for item in session_runs: + result.append(item) + return validated_json(apispec.SessionRuns, result) + + return "/persisted_logs/sessions//runs", ["GET"], _get_session_runs + + def get_build_logs(self) -> BlueprintFactoryResponse: + """Get persisted image build logs.""" + + @authenticate(self.authenticator) + @only_authenticated + async def _get_build_logs(_: Request, user: base_models.APIUser, build_id: ULID) -> JSONResponse: + async with self.session_maker() as session, session.begin(): + result = await self.build_logs_repo.get_build_logs(session=session, user=user, build_id=build_id) + return validated_json(apispec.PersistedBuildLogs, dict(logs=result)) + + return "/persisted_logs/builds/", ["GET"], _get_build_logs diff --git a/components/renku_data_services/persisted_logs/collector.py b/components/renku_data_services/persisted_logs/collector.py new file mode 100644 index 0000000000..83dfcf353c --- /dev/null +++ b/components/renku_data_services/persisted_logs/collector.py @@ -0,0 +1,334 @@ +"""Collector for gathering persisted logs.""" + +import asyncio +from abc import abstractmethod +from collections.abc import AsyncIterator, Callable +from datetime import UTC, datetime, timedelta + +import httpx +from pydantic import ValidationError +from sqlalchemy.ext.asyncio import AsyncSession +from ulid import ULID + +from renku_data_services.app_config import logging +from renku_data_services.persisted_logs import loki_api, models +from renku_data_services.persisted_logs.config import PersistedLogsConfig +from renku_data_services.persisted_logs.constants import ( + ONE_MINUTE_IN_NANOS, + PERSISTED_LOGS_BUILD_LABEL_KEY, + PERSISTED_LOGS_BUILD_LABEL_VALUE, + PERSISTED_LOGS_NAMESPACE_LABEL_KEY, + PERSISTED_LOGS_SESSIONS_LABEL_KEY, + PERSISTED_LOGS_SESSIONS_LABEL_VALUE, +) +from renku_data_services.persisted_logs.db import ( + AmaltheaSessionPersistedLogsWriteRepository, + ImageBuildPersistedLogsWriteRepository, +) + +logger = logging.getLogger(__name__) + + +class LokiLogReader: + """Read logs from loki.""" + + def __init__(self, config: PersistedLogsConfig, client: httpx.AsyncClient) -> None: + self.config = config + self.client = client + self.client.base_url = httpx.URL(config.loki_read_base_url) + + async def get_amalthea_session_logs( + self, limit: int = 1000, start: int | None = None, end: int | None = None + ) -> AsyncIterator[models.UnsavedSessionLogLine]: + """Fetches Amalthea session logs from Loki. + + Parameters: + - limit: max number of entries to return + - start: start timestamp as a Unix nano timestamp + - end: end timestamp as a Unix nano timestamp + + See also https://grafana.com/docs/loki/latest/reference/loki-http-api/#query-logs-within-a-range-of-time + """ + query = ( + "{" + f'{PERSISTED_LOGS_SESSIONS_LABEL_KEY}="{PERSISTED_LOGS_SESSIONS_LABEL_VALUE}",' + f'{PERSISTED_LOGS_NAMESPACE_LABEL_KEY}="{self.config.namespace}"' + "}" + ) + response = await self._get_logs(query=query, limit=limit, start=start, end=end) + async for item in self._process_session_logs(response): + yield item + + async def get_image_build_logs( + self, limit: int = 1000, start: int | None = None, end: int | None = None + ) -> AsyncIterator[models.UnsavedBuildLogLine]: + """Fetchesimage build logs from Loki. + + Parameters: + - limit: max number of entries to return + - start: start timestamp as a Unix nano timestamp + - end: end timestamp as a Unix nano timestamp + + See also https://grafana.com/docs/loki/latest/reference/loki-http-api/#query-logs-within-a-range-of-time + """ + query = ( + "{" + f'{PERSISTED_LOGS_BUILD_LABEL_KEY}="{PERSISTED_LOGS_BUILD_LABEL_VALUE}",' + f'{PERSISTED_LOGS_NAMESPACE_LABEL_KEY}="{self.config.namespace}"' + "}" + ) + response = await self._get_logs(query=query, limit=limit, start=start, end=end) + async for item in self._process_image_build_logs(response): + yield item + + async def _get_logs( + self, query: str, limit: int = 1000, start: int | None = None, end: int | None = None + ) -> loki_api.LokiQueryRangeResponse: + """Fetches logs from Loki, using the passed in query. + + Parameters: + - query: the Loki query + - limit: max number of entries to return + - start: start timestamp as a Unix nano timestamp + - end: end timestamp as a Unix nano timestamp + + See also https://grafana.com/docs/loki/latest/reference/loki-http-api/#query-logs-within-a-range-of-time + """ + params: dict[str, str | int] = dict() + params["query"] = query + params["direction"] = "forward" + params["limit"] = limit + if start: + params["start"] = str(start) + if end: + params["end"] = str(end) + res = await self.client.get("loki/api/v1/query_range", params=params) + res.raise_for_status() + return loki_api.LokiQueryRangeResponse.model_validate_json(res.content) + + @staticmethod + async def _process_session_logs( + response: loki_api.LokiQueryRangeResponse, + ) -> AsyncIterator[models.UnsavedSessionLogLine]: + log_line_ids: set[str] = set() + for entry in response.data.result: + stream: loki_api.AmaltheaSessionStream | None = None + try: + stream = loki_api.AmaltheaSessionStream.model_validate(entry.stream) + except ValidationError as err: + logger.warning(f"Skipping entry {entry.stream} because of validation error: {err}") + continue + + try: + launcher_id = ULID.from_str(stream.renku_io_launcher_id.upper()) + except ValueError as err: + logger.warning( + f"Skipping entry {entry.stream} because renku_io_launcher_id='{stream.renku_io_launcher_id}' " + f"is not a valid ULID: {err}" + ) + continue + + try: + run_id = ULID.from_str(stream.renku_io_run_id.upper()) + except ValueError as err: + logger.warning( + f"Skipping entry {entry.stream} because renku_io_run_id='{stream.renku_io_run_id}' " + f"is not a valid ULID: {err}" + ) + continue + + for nano_ts, log_line in entry.values: + log_line_id = f"{nano_ts.root}::{stream.container}::{stream.pod}" + + if log_line_id in log_line_ids: + continue + + log_line_ids.add(log_line_id) + yield models.UnsavedSessionLogLine( + id=log_line_id, + user_id=stream.renku_io_safe_username, + run_id=run_id, + session_uid=stream.renku_io_session_uid, + launcher_id=launcher_id, + submission_id=stream.renku_io_submission_id, + container=stream.container, + timestamp=nano_ts.get_value(), + log_line=log_line, + ) + + @staticmethod + async def _process_image_build_logs( + response: loki_api.LokiQueryRangeResponse, + ) -> AsyncIterator[models.UnsavedBuildLogLine]: + log_line_ids: set[str] = set() + for entry in response.data.result: + stream: loki_api.ShipwrightBuildRunStream | None = None + try: + stream = loki_api.ShipwrightBuildRunStream.model_validate(entry.stream) + except ValidationError as err: + logger.warning(f"Skipping entry {entry.stream} because of validation error: {err}") + continue + + try: + build_id = ULID.from_str(stream.renku_io_buildrun_name.upper().removeprefix("RENKU-")) + except ValueError as err: + logger.warning( + f"Skipping entry {entry.stream} because renku_io_buildrun_name='{stream.renku_io_buildrun_name}' " + f"does not contain a valid ULID: {err}" + ) + continue + + for nano_ts, log_line in entry.values: + log_line_id = f"{nano_ts.root}::{stream.container}::{stream.pod}" + + if log_line_id in log_line_ids: + continue + + log_line_ids.add(log_line_id) + yield models.UnsavedBuildLogLine( + id=log_line_id, + build_id=build_id, + container=stream.container, + timestamp=nano_ts.get_value(), + log_line=log_line, + ) + + +class PersistedLogsCollector: + """Abstract class for gathering persisted logs.""" + + @abstractmethod + async def collect_persisted_logs(self) -> None: + """Collect persisted logs from Amalthea sessions and image builds.""" + ... + + @abstractmethod + async def purge_expired_logs(self) -> None: + """Purge expired persisted logs from the database.""" + ... + + @staticmethod + def from_config( + config: PersistedLogsConfig, + session_maker: Callable[..., AsyncSession], + http_client: httpx.AsyncClient | None = None, + ) -> "PersistedLogsCollector": + """Construct a PersistedLogsCollector from a configuration object.""" + if config.enabled: + if http_client is None: + http_client = httpx.AsyncClient() + reader = LokiLogReader(config=config, client=http_client) + return DefaultPersistedLogsCollector( + session_maker=session_maker, + config=config, + reader=reader, + session_logs_repo=AmaltheaSessionPersistedLogsWriteRepository(), + build_logs_repo=ImageBuildPersistedLogsWriteRepository(), + ) + return NoopPersistedLogsCollector() + + +class NoopPersistedLogsCollector(PersistedLogsCollector): + """No-op collector.""" + + async def collect_persisted_logs(self) -> None: + """Collect persisted logs from Amalthea sessions and image builds.""" + return None + + async def purge_expired_logs(self) -> None: + """Purge expired persisted logs from the database.""" + return None + + +class DefaultPersistedLogsCollector(PersistedLogsCollector): + """Collector for gathering persisted logs.""" + + def __init__( + self, + session_maker: Callable[..., AsyncSession], + config: PersistedLogsConfig, + reader: LokiLogReader, + session_logs_repo: AmaltheaSessionPersistedLogsWriteRepository, + build_logs_repo: ImageBuildPersistedLogsWriteRepository, + ) -> None: + self.session_maker = session_maker + self.config = config + self.reader = reader + self.session_logs_repo = session_logs_repo + self.build_logs_repo = build_logs_repo + + async def collect_persisted_logs(self) -> None: + """Collect persisted logs from Amalthea sessions and image builds.""" + await asyncio.gather(self.collect_sessions_persisted_logs(), self.collect_build_persisted_logs()) + return None + + async def collect_sessions_persisted_logs(self) -> None: + """Collect persisted logs from Amalthea sessions.""" + + async with self.session_maker() as session: + async with session.begin(): + ts = await self.session_logs_repo.get_latest_log_timestamp(session=session) + start = _one_hour_ago_in_nanos() + if ts is not None and ts > start: + start = ts - ONE_MINUTE_IN_NANOS + + # Loop to collect all logs, including late entries + has_more = True + current_start = start + while has_more: + logs_stream = self.reader.get_amalthea_session_logs(start=current_start) + async with session.begin(): + result = await self.session_logs_repo.insert_session_logs(session=session, logs_stream=logs_stream) + current_start = result.last_timestamp + 1 + has_more = result.log_count > 1 + + return None + + async def collect_build_persisted_logs(self) -> None: + """Collect persisted logs from image builds.""" + + async with self.session_maker() as session: + async with session.begin(): + ts = await self.build_logs_repo.get_latest_log_timestamp(session=session) + start = _one_hour_ago_in_nanos() + if ts is not None and ts > start: + start = ts - ONE_MINUTE_IN_NANOS + + # Loop to collect all logs, including late entries + has_more = True + current_start = start + while has_more: + logs_stream = self.reader.get_image_build_logs(start=current_start) + async with session.begin(): + result = await self.build_logs_repo.insert_build_logs(session=session, logs_stream=logs_stream) + current_start = result.last_timestamp + 1 + has_more = result.log_count > 1 + + return None + + async def purge_expired_logs(self) -> None: + """Purge expired persisted logs from the database.""" + await asyncio.gather(self.purge_expired_session_logs(), self.purge_expired_build_logs()) + return None + + async def purge_expired_session_logs(self) -> None: + """Purge expired session logs from the database.""" + now = datetime.now(tz=UTC) + cutoff = now - self.config.logs_ttl + async with self.session_maker() as session, session.begin(): + await self.session_logs_repo.delete_expired_session_logs(session=session, before=cutoff) + return None + + async def purge_expired_build_logs(self) -> None: + """Purge expired image build logs from the database.""" + now = datetime.now(tz=UTC) + cutoff = now - self.config.logs_ttl + async with self.session_maker() as session, session.begin(): + await self.build_logs_repo.delete_expired_build_logs(session=session, before=cutoff) + return None + + +def _one_hour_ago_in_nanos() -> int: + """Returns the Unix nano timestamp corresponding to one hour ago.""" + dt = datetime.now(tz=UTC) - timedelta(hours=1) + return int(dt.timestamp() * 1e6) * 1000 diff --git a/components/renku_data_services/persisted_logs/config.py b/components/renku_data_services/persisted_logs/config.py new file mode 100644 index 0000000000..bde4897c4f --- /dev/null +++ b/components/renku_data_services/persisted_logs/config.py @@ -0,0 +1,29 @@ +"""Configuration for persisted logs.""" + +from dataclasses import dataclass +from datetime import timedelta + + +@dataclass(eq=True, frozen=True, kw_only=True) +class PersistedLogsConfig: + """Configuration for persisted logs.""" + + enabled: bool + loki_read_base_url: str + namespace: str + logs_ttl: timedelta + + @classmethod + def from_env(cls, namespace: str) -> "PersistedLogsConfig": + """Create a config from environment variables.""" + # enabled = os.environ.get("PERSISTED_LOG_ENABLED", "false").lower() == "true" + # return cls( + # enabled=enabled, + # ) + # TODO: load config from env vars + return cls( + enabled=True, + loki_read_base_url="http://loki-read.monitoring.svc.cluster.local:3100/", + namespace=namespace, + logs_ttl=timedelta(days=1), + ) diff --git a/components/renku_data_services/persisted_logs/constants.py b/components/renku_data_services/persisted_logs/constants.py new file mode 100644 index 0000000000..e55b44201d --- /dev/null +++ b/components/renku_data_services/persisted_logs/constants.py @@ -0,0 +1,30 @@ +"""Constants for persisted logs.""" + +from typing import Final + +PERSISTED_LOGS_SESSIONS_LABEL_KEY: Final[str] = "app" +"""The loki label key to select session logs streams.""" + +PERSISTED_LOGS_SESSIONS_LABEL_VALUE: Final[str] = "AmaltheaSession" +"""The loki label value to select session logs streams.""" + +PERSISTED_LOGS_BUILD_LABEL_KEY: Final[str] = "app" +"""The loki label key to select build logs streams.""" + +PERSISTED_LOGS_BUILD_LABEL_VALUE: Final[str] = "ShipwrightBuildRun" +"""The loki label value to select build logs streams.""" + +PERSISTED_LOGS_NAMESPACE_LABEL_KEY: Final[str] = "namespace" +"""The loki label key to select logs streams from a specific kubernetes namespace.""" + +ONE_SECOND_IN_NANOS: Final[int] = 1_000_000_000 +"""One second as nanoseconds (for Loki).""" + +ONE_MINUTE_IN_NANOS: Final[int] = 60 * ONE_SECOND_IN_NANOS +"""One minute as nanoseconds (for Loki).""" + +SESSION_MAIN_CONTAINER: Final[str] = "amalthea-session" +"""The name of the main pod container for Amalthea sessions.""" + +BUILD_MAIN_CONTAINER: Final[str] = "step-build-and-push" +"""The name of the main pod container for image builds.""" diff --git a/components/renku_data_services/persisted_logs/db.py b/components/renku_data_services/persisted_logs/db.py new file mode 100644 index 0000000000..263536401a --- /dev/null +++ b/components/renku_data_services/persisted_logs/db.py @@ -0,0 +1,421 @@ +"""Adapters for persisted logs database classes.""" + +from __future__ import annotations + +from collections.abc import AsyncIterator, Sequence +from datetime import datetime +from typing import TYPE_CHECKING + +from sqlalchemy import delete, select +from sqlalchemy.exc import DatabaseError +from sqlalchemy.ext.asyncio import AsyncScalarResult, AsyncSession +from ulid import ULID + +from renku_data_services import base_models, errors +from renku_data_services.app_config import logging +from renku_data_services.authz.authz import Authz, ResourceType +from renku_data_services.authz.models import Scope +from renku_data_services.persisted_logs import models +from renku_data_services.persisted_logs import orm as schemas +from renku_data_services.persisted_logs.constants import BUILD_MAIN_CONTAINER, SESSION_MAIN_CONTAINER +from renku_data_services.repositories import models as repo_models +from renku_data_services.session import models as session_models +from renku_data_services.session import orm as session_schemas + +if TYPE_CHECKING: + from renku_data_services.repositories.db import GitRepositoriesRepository + from renku_data_services.session.config import BuildsConfig +logger = logging.getLogger(__name__) + + +class AmaltheaSessionPersistedLogsReadRepository: + """Repository for reading persisted logs of Amalthea sessions.""" + + def __init__(self, authz: Authz) -> None: + self.authz: Authz = authz + + async def get_session_logs( + self, + session: AsyncSession, + user: base_models.APIUser, + launcher_id: ULID, + run_id: ULID | None = None, + submission_id: str | None = None, + ) -> models.PersistedSessionLogs | None: + """Returns persisted session logs for the given launcher.""" + if not user.is_authenticated or not user.id: + raise errors.UnauthorizedError(message="You have to be authenticated to perform this operation.") + await self._check_session_launcher(session=session, user=user, launcher_id=launcher_id) + session_run = await self._get_session_run( + session=session, user_id=user.id, launcher_id=launcher_id, run_id=run_id, submission_id=submission_id + ) + if session_run is None: + return None + logs_per_container = await self._get_logs_per_container(session=session, run_id=session_run.id) + return models.PersistedSessionLogs( + run=session_run, + logs=logs_per_container, + ) + + async def get_session_runs( + self, + session: AsyncSession, + user: base_models.APIUser, + launcher_id: ULID, + ) -> AsyncIterator[models.SessionRun]: + """Returns the session runs for the given launcher.""" + if not user.is_authenticated or not user.id: + raise errors.UnauthorizedError(message="You have to be authenticated to perform this operation.") + await self._check_session_launcher(session=session, user=user, launcher_id=launcher_id) + stmt = ( + select(schemas.SessionRunORM) + .where(schemas.SessionRunORM.user_id == user.id) + .where(schemas.SessionRunORM.launcher_id == launcher_id) + .order_by(schemas.SessionRunORM.id.desc()) + ) + res = await session.stream_scalars(stmt) + async for session_run_orm in res: + yield session_run_orm.dump() + + async def _check_session_launcher( + self, session: AsyncSession, user: base_models.APIUser, launcher_id: ULID + ) -> None: + """Check that the session launcher exists and the user has access to it.""" + stmt = select(session_schemas.SessionLauncherORM).where(session_schemas.SessionLauncherORM.id == launcher_id) + res = await session.scalars(stmt) + launcher_orm = res.one_or_none() + authorized = ( + await self.authz.has_permission(user, ResourceType.project, launcher_orm.project_id, Scope.READ) + if launcher_orm is not None + else False + ) + if not authorized or launcher_orm is None: + raise errors.MissingResourceError( + message=f"Session launcher with id '{launcher_id}' does not exist or you do not have access to it." + ) + + async def _get_session_run( + self, + session: AsyncSession, + user_id: str, + launcher_id: ULID, + run_id: ULID | None = None, + submission_id: str | None = None, + ) -> models.SessionRun | None: + """Get a specific session run from the persisted logs database. + + If no `run_id` is specified, then return the latest session run. + """ + stmt = ( + select(schemas.SessionRunORM) + .where(schemas.SessionRunORM.user_id == user_id) + .where(schemas.SessionRunORM.launcher_id == launcher_id) + .order_by(schemas.SessionRunORM.id.desc()) + .limit(1) + ) + if run_id: + stmt = stmt.where(schemas.SessionRunORM.id == run_id) + if submission_id: + stmt = stmt.where(schemas.SessionRunORM.submission_id == submission_id) + res = await session.scalars(stmt) + session_run_orm = res.one_or_none() + if session_run_orm is None: + return None + return session_run_orm.dump() + + async def _get_logs_per_container(self, session: AsyncSession, run_id: ULID) -> Sequence[models.ContainerLogs]: + """Get the logs of a specific session run, organized by container.""" + # TODO: handle pagination? + stmt = ( + select(schemas.AmaltheaSessionLogORM) + .where(schemas.AmaltheaSessionLogORM.run_id == run_id) + .order_by(schemas.AmaltheaSessionLogORM.id.asc()) + ) + res = await session.stream_scalars(stmt) + # Sort logs by container name, forcing "amalthea-session" to be the first item (main container) + return await _sort_logs_per_container(res, main_container=SESSION_MAIN_CONTAINER) + + +class AmaltheaSessionPersistedLogsWriteRepository: + """Repository for writing persisted logs of Amalthea sessions. + + The write side is performed as a background task and does not access authz. + """ + + async def get_latest_log_timestamp(self, session: AsyncSession) -> int | None: + """Returns the latest log timestamp.""" + stmt = ( + select(schemas.AmaltheaSessionLogORM.timestamp) + .select_from(schemas.AmaltheaSessionLogORM) + .order_by(schemas.AmaltheaSessionLogORM.timestamp.desc()) + .limit(1) + ) + res = await session.scalars(stmt) + timestamp = res.one_or_none() + return timestamp + + async def insert_session_logs( + self, session: AsyncSession, logs_stream: AsyncIterator[models.UnsavedSessionLogLine] + ) -> models.LogStreamMetadata: + """Insert sessions logs into the persisted logs database.""" + log_count = 0 + last_timestamp = 0 + async for log in logs_stream: + log_count += 1 + if log.timestamp > last_timestamp: + last_timestamp = log.timestamp + try: + await self._insert_log_line(session=session, log=log) + except DatabaseError as err: + logger.warning(f"Could not process log line {log.id}: {err}") + + return models.LogStreamMetadata(log_count=log_count, last_timestamp=last_timestamp) + + async def delete_expired_session_logs(self, session: AsyncSession, before: datetime) -> int: + """Remove expired session logs from the database.""" + nano_ts = models.NanoTimestamp.from_datetime(before) + delete_logs_stmt = delete(schemas.AmaltheaSessionLogORM).where( + schemas.AmaltheaSessionLogORM.timestamp < nano_ts + ) + res = await session.execute(delete_logs_stmt) + deleted_logs_count = res.rowcount + + # Remove orphaned session runs + stmt = ( + select(schemas.SessionRunORM.id) + .join( + schemas.AmaltheaSessionLogORM, + schemas.SessionRunORM.id == schemas.AmaltheaSessionLogORM.run_id, + isouter=True, # isouter makes it a left-join, not an outer join + ) + .where(schemas.AmaltheaSessionLogORM.id.is_(None)) + ) + session_runs_res = await session.scalars(stmt) + session_run_ids = session_runs_res.all() + await session.execute(delete(schemas.SessionRunORM).where(schemas.SessionRunORM.id.in_(session_run_ids))) + + return deleted_logs_count + + async def _insert_log_line(self, session: AsyncSession, log: models.UnsavedSessionLogLine) -> bool: + """Insert a single session log line into the persisted logs database. + + Returns true if the log line was inserted into the database and false otherwise (the log line already exists). + """ + existing_log_res = await session.scalars( + select(schemas.AmaltheaSessionLogORM.id).where(schemas.AmaltheaSessionLogORM.id == log.id) + ) + existing_log_orm = existing_log_res.one_or_none() + if existing_log_orm: + return False + + session_run_res = await session.scalars( + select(schemas.SessionRunORM).where(schemas.SessionRunORM.id == log.run_id) + ) + session_run_orm = session_run_res.one_or_none() + if session_run_orm is None: + async with session.begin_nested(): + session_run_orm = schemas.SessionRunORM( + id=log.run_id, + user_id=log.user_id, + session_uid=log.session_uid, + launcher_id=log.launcher_id, + submission_id=log.submission_id, + ) + session.add(session_run_orm) + await session.flush() + + log_orm = schemas.AmaltheaSessionLogORM( + id=log.id, + run_id=log.run_id, + container=log.container, + timestamp=log.timestamp, + log_line=log.log_line, + ) + session.add(log_orm) + await session.flush() + return True + + +class ImageBuildPersistedLogsReadRepository: + """Repository for reading persisted logs of image builds.""" + + def __init__( + self, + authz: Authz, + builds_config: BuildsConfig, + git_repositories_repo: GitRepositoriesRepository, + ) -> None: + self.authz: Authz = authz + self.builds_config = builds_config + self.git_repositories_repo = git_repositories_repo + + async def get_build_logs( + self, session: AsyncSession, user: base_models.APIUser, build_id: ULID + ) -> Sequence[models.ContainerLogs]: + """Returns persisted session logs for the given image build.""" + if not user.is_authenticated or not user.id: + raise errors.UnauthorizedError(message="You have to be authenticated to perform this operation.") + await self._check_build(session=session, user=user, build_id=build_id) + logs_per_container = await self._get_logs_per_container(session=session, build_id=build_id) + return logs_per_container + + async def _check_build(self, session: AsyncSession, user: base_models.APIUser, build_id: ULID) -> None: + """Check that the image build exists and the user has access to it.""" + stmt = select(session_schemas.BuildORM).where(session_schemas.BuildORM.id == build_id) + res = await session.scalars(stmt) + build_orm = res.one_or_none() + authorized = ( + await self._check_environment( + session=session, user=user, environment=build_orm.environment, scope=Scope.READ + ) + if build_orm is not None + else False + ) + + # If the output image is private, check that the user can read the source repository + if build_orm is None or build_orm.result_image is None: + authorized = False + else: + if self.builds_config.private_builds_enabled and build_orm.result_image.startswith( + self.builds_config.build_output_private_image_prefix + ): + if build_orm.result_repository_url is None: + authorized = False + else: + repo_data = await self.git_repositories_repo.get_repository( + repository_url=build_orm.result_repository_url, + user=user, + etag=None, + ) + if ( + not isinstance(repo_data.metadata, repo_models.Metadata) + or not repo_data.metadata.pull_permission + ): + authorized = False + + if not authorized or build_orm is None: + raise errors.MissingResourceError( + message=f"Build with id '{build_id}' does not exist or you do not have access to it." + ) + + async def _check_environment( + self, + session: AsyncSession, + user: base_models.APIUser, + environment: session_schemas.EnvironmentORM, + scope: Scope, + ) -> bool: + """Checks whether the provided user has a specific permission on a session environment.""" + if environment.environment_kind == session_models.EnvironmentKind.GLOBAL: + return scope == Scope.READ or user.is_admin + + launcher = await session.scalar( + select(schemas.SessionLauncherORM).where(schemas.SessionLauncherORM.environment_id == environment.id) + ) + authorized = False + if launcher: + authorized = await self.authz.has_permission(user, ResourceType.project, launcher.project_id, scope) + return authorized + + async def _get_logs_per_container(self, session: AsyncSession, build_id: ULID) -> Sequence[models.ContainerLogs]: + """Get the logs of a specific image build, organized by container.""" + # TODO: handle pagination? + stmt = ( + select(schemas.ImageBuildLogORM) + .where(schemas.ImageBuildLogORM.build_id == build_id) + .order_by(schemas.ImageBuildLogORM.id.asc()) + ) + res = await session.stream_scalars(stmt) + # Sort logs by container name, forcing "step-build-and-push" to be the first item (main container) + return await _sort_logs_per_container(res, main_container=BUILD_MAIN_CONTAINER) + + +class ImageBuildPersistedLogsWriteRepository: + """Repository for writing persisted logs of image builds. + + The write side is performed as a background task and does not access authz. + """ + + async def get_latest_log_timestamp(self, session: AsyncSession) -> int | None: + """Returns the latest log timestamp.""" + stmt = ( + select(schemas.ImageBuildLogORM.timestamp) + .select_from(schemas.ImageBuildLogORM) + .order_by(schemas.ImageBuildLogORM.timestamp.desc()) + .limit(1) + ) + res = await session.scalars(stmt) + timestamp = res.one_or_none() + return timestamp + + async def insert_build_logs( + self, session: AsyncSession, logs_stream: AsyncIterator[models.UnsavedBuildLogLine] + ) -> models.LogStreamMetadata: + """Insert sessions logs into the persisted logs database.""" + log_count = 0 + last_timestamp = 0 + async for log in logs_stream: + log_count += 1 + if log.timestamp > last_timestamp: + last_timestamp = log.timestamp + try: + await self._insert_log_line(session=session, log=log) + except DatabaseError as err: + logger.warning(f"Could not process log line {log.id}: {err}") + + return models.LogStreamMetadata(log_count=log_count, last_timestamp=last_timestamp) + + async def delete_expired_build_logs(self, session: AsyncSession, before: datetime) -> int: + """Remove expired build logs from the database.""" + nano_ts = models.NanoTimestamp.from_datetime(before) + delete_logs_stmt = delete(schemas.ImageBuildLogORM).where(schemas.ImageBuildLogORM.timestamp < nano_ts) + res = await session.execute(delete_logs_stmt) + deleted_logs_count = res.rowcount + return deleted_logs_count + + async def _insert_log_line(self, session: AsyncSession, log: models.UnsavedBuildLogLine) -> bool: + """Insert a single session log line into the persisted logs database. + + Returns true if the log line was inserted into the database and false otherwise (the log line already exists). + """ + existing_log_res = await session.scalars( + select(schemas.ImageBuildLogORM.id).where(schemas.ImageBuildLogORM.id == log.id) + ) + existing_log_orm = existing_log_res.one_or_none() + if existing_log_orm: + return False + + async with session.begin_nested(): + log_orm = schemas.ImageBuildLogORM( + id=log.id, + build_id=log.build_id, + container=log.container, + timestamp=log.timestamp, + log_line=log.log_line, + ) + session.add(log_orm) + await session.flush() + return True + + +async def _sort_logs_per_container( + result: AsyncScalarResult[schemas.AmaltheaSessionLogORM] | AsyncScalarResult[schemas.ImageBuildLogORM], + main_container: str | None = None, +) -> Sequence[models.ContainerLogs]: + """Organize logs per container.""" + logs_per_container: dict[str, list[models.LogLine]] = dict() + async for log_entry in result: + container = log_entry.container + logs = logs_per_container.get(container) + if logs is None: + logs = list[models.LogLine]() + logs_per_container[container] = logs + logs.append(models.LogLine(timestamp=log_entry.timestamp, log_line=log_entry.log_line)) + # Sort containers by name, forcing `main_container` to be the first item + containers_set = set(logs_per_container.keys()) + containers: list[str] = list() + if main_container and main_container in containers_set: + containers.append(main_container) + containers_set.remove(main_container) + containers.extend(sorted(containers_set)) + return [models.ContainerLogs(container=container, logs=logs_per_container[container]) for container in containers] diff --git a/components/renku_data_services/persisted_logs/loki_api.py b/components/renku_data_services/persisted_logs/loki_api.py new file mode 100644 index 0000000000..3ce17eda9e --- /dev/null +++ b/components/renku_data_services/persisted_logs/loki_api.py @@ -0,0 +1,81 @@ +"""Pydantic models for the Loki API.""" + +from __future__ import annotations + +from enum import StrEnum + +from pydantic import BaseModel, ConfigDict, Field, RootModel + + +class Base(BaseModel): + """Base CRD specification.""" + + model_config = ConfigDict( + # Do not exclude unknown properties. + extra="allow" + ) + + +class LokiQueryRangeResponse(Base): + """Response from the query range endpoint (streams only).""" + + status: LokiQueryRangeResponseStatus + data: LokiQueryRangeResponseData + + +class LokiQueryRangeResponseStatus(StrEnum): + """Response status.""" + + success = "success" + + +class LokiQueryRangeResponseData(Base): + """Response data from the query range endpoint (streams only).""" + + result_type: LokiQueryRangeResponseResultType = Field(..., alias="resultType") + result: list[LokiQueryRangeResponseStream] + + +class LokiQueryRangeResponseResultType(StrEnum): + """Result type.""" + + streams = "streams" + + +class LokiQueryRangeResponseStream(Base): + """Loki log stream.""" + + stream: dict[str, str] + values: list[tuple[NanoTimestamp, str]] + + +class NanoTimestamp(RootModel[str]): + """Unix timestamp in nanoseconds.""" + + root: str = Field(..., pattern="\\d+") + + def get_value(self) -> int: + """Return the timestamp as a big integer.""" + return int(self.root) + + +class AmaltheaSessionStream(Base): + """Loki stream labels for logs extracted from an Amalthea session.""" + + container: str + pod: str + renku_io_launcher_id: str + renku_io_project_id: str | None = None + renku_io_run_id: str + renku_io_safe_username: str + renku_io_session_type: str | None = None + renku_io_session_uid: str | None = None + renku_io_submission_id: str | None = None + + +class ShipwrightBuildRunStream(Base): + """Loki stream labels for logs extracted from a Shipwright build run.""" + + container: str + pod: str + renku_io_buildrun_name: str diff --git a/components/renku_data_services/persisted_logs/models.py b/components/renku_data_services/persisted_logs/models.py new file mode 100644 index 0000000000..ea48a60937 --- /dev/null +++ b/components/renku_data_services/persisted_logs/models.py @@ -0,0 +1,102 @@ +"""Models for persisted logs.""" + +from collections.abc import Sequence +from dataclasses import dataclass +from datetime import UTC, datetime +from typing import Self + +from ulid import ULID + + +class NanoTimestamp(int): + """Unix timestamp in nanoseconds.""" + + def to_datetime(self) -> datetime: + """Return the corresponding datetime, trucated to microsecond precision.""" + return datetime.fromtimestamp((self // 1_000) / 1e6, tz=UTC) + + @classmethod + def from_datetime(cls, dt: datetime) -> Self: + """Create a nano timestamp from a datetime object.""" + return cls(int(dt.timestamp() * 1e6) * 1000) + + +@dataclass(eq=True, frozen=True, kw_only=True) +class UnsavedSessionLogLine: + """Represents an unsaved log line.""" + + id: str + """The ID of the log line. + + This is used to de-duplicate log lines. + """ + + user_id: str + run_id: ULID + session_uid: str | None + launcher_id: ULID + submission_id: str | None + container: str + timestamp: int + log_line: str + + +@dataclass(eq=True, frozen=True, kw_only=True) +class SessionRun: + """The continuous execution span of a session.""" + + id: ULID + session_uid: str | None + launcher_id: ULID + submission_id: str | None + + +@dataclass(eq=True, frozen=True, kw_only=True) +class UnsavedBuildLogLine: + """Represents an unsaved image build log line.""" + + id: str + """The ID of the log line. + + This is used to de-duplicate log lines. + """ + + build_id: ULID + container: str + timestamp: int + log_line: str + + +@dataclass(eq=True, frozen=True, kw_only=True) +class LogLine: + """A single log line.""" + + timestamp: int + log_line: str + + +@dataclass(eq=True, frozen=True, kw_only=True) +class ContainerLogs: + """Logs of a single container.""" + + container: str + logs: Sequence[LogLine] + + +@dataclass(eq=True, frozen=True, kw_only=True) +class PersistedSessionLogs: + """Result of getting session logs from the database.""" + + run: SessionRun + logs: Sequence[ContainerLogs] + + +@dataclass(eq=True, frozen=True, kw_only=True) +class LogStreamMetadata: + """Log stream metadata. + + Used to know if there are more logs to fetch and if so, where to continue from. + """ + + log_count: int + last_timestamp: int diff --git a/components/renku_data_services/persisted_logs/orm.py b/components/renku_data_services/persisted_logs/orm.py new file mode 100644 index 0000000000..9c6476d5fe --- /dev/null +++ b/components/renku_data_services/persisted_logs/orm.py @@ -0,0 +1,95 @@ +"""SQLAlchemy schemas for the peristed logs database.""" + +from __future__ import annotations + +from sqlalchemy import BigInteger, ForeignKey, MetaData +from sqlalchemy.orm import DeclarativeBase, Mapped, MappedAsDataclass, mapped_column, relationship +from ulid import ULID + +from renku_data_services.base_orm.registry import COMMON_ORM_REGISTRY +from renku_data_services.persisted_logs import models +from renku_data_services.session.orm import BuildORM, SessionLauncherORM +from renku_data_services.users.orm import UserORM +from renku_data_services.utils.sqlalchemy import ULIDType + + +class BaseORM(MappedAsDataclass, DeclarativeBase): + """Base class for all ORM classes.""" + + metadata = MetaData(schema="persisted_logs") + registry = COMMON_ORM_REGISTRY + + +class SessionRunORM(BaseORM): + """A session run, which is the continuous execution of a session.""" + + __tablename__ = "session_runs" + + id: Mapped[ULID] = mapped_column("id", ULIDType, primary_key=True) + """ID of a session run.""" + + user_id: Mapped[str] = mapped_column(ForeignKey(UserORM.keycloak_id), index=True, nullable=False) + """User ID of the owner of the session.""" + + session_uid: Mapped[str | None] = mapped_column(nullable=True) + """The session UID for this session run.""" + + launcher_id: Mapped[ULID] = mapped_column(ULIDType, ForeignKey(SessionLauncherORM.id), index=True, nullable=False) + """The session launcher ID of the session.""" + + submission_id: Mapped[str | None] = mapped_column(nullable=True) + """The submission ID, if the session run corresponds to an offline job.""" + + def dump(self) -> models.SessionRun: + """Create a session run model from the SessionRunORM.""" + return models.SessionRun( + id=self.id, + session_uid=self.session_uid, + launcher_id=self.launcher_id, + submission_id=self.submission_id, + ) + + +class AmaltheaSessionLogORM(BaseORM): + """A log line from an Amalthea session.""" + + __tablename__ = "amalthea_session_logs" + + id: Mapped[str] = mapped_column("id", primary_key=True, nullable=False) + """ID of the log line.""" + + run_id: Mapped[ULID] = mapped_column(ForeignKey(SessionRunORM.id, ondelete="CASCADE"), index=True, nullable=False) + """ID of the session run.""" + + session_run: Mapped[SessionRunORM] = relationship(lazy="select", init=False, repr=False, viewonly=True) + """The session run this log line belongs to.""" + + container: Mapped[str] = mapped_column(nullable=False) + """The container this log line belongs to.""" + + timestamp: Mapped[int] = mapped_column(BigInteger, index=True, nullable=False) + """The timestamp of the log line (nanosecond timestamp).""" + + log_line: Mapped[str] = mapped_column(nullable=False) + """The contents of the log line.""" + + +class ImageBuildLogORM(BaseORM): + """A log line from an image build.""" + + __tablename__ = "image_build_logs" + + id: Mapped[str] = mapped_column("id", primary_key=True, nullable=False) + """ID of the log line.""" + + build_id: Mapped[ULID] = mapped_column(ForeignKey(BuildORM.id, ondelete="CASCADE"), index=True, nullable=False) + """ID of the image build.""" + + container: Mapped[str] = mapped_column(nullable=False) + """The container this log line belongs to.""" + + timestamp: Mapped[int] = mapped_column(BigInteger, index=True, nullable=False) + """The timestamp of the log line (nanosecond timestamp).""" + + log_line: Mapped[str] = mapped_column(nullable=False) + """The contents of the log line.""" diff --git a/projects/renku_data_service/pyproject.toml b/projects/renku_data_service/pyproject.toml index 0341633576..db85b6d591 100644 --- a/projects/renku_data_service/pyproject.toml +++ b/projects/renku_data_service/pyproject.toml @@ -44,6 +44,7 @@ packages = [ { include = "renku_data_services/metrics", from = "../../components" }, { include = "renku_data_services/capacity_reservation", from = "../../components" }, { include = "renku_data_services/resource_usage", from = "../../components" }, + { include = "renku_data_services/persisted_logs", from = "../../components" }, ] [tool.poetry.dependencies] diff --git a/projects/renku_data_tasks/pyproject.toml b/projects/renku_data_tasks/pyproject.toml index d816a42841..9a5dc783f4 100644 --- a/projects/renku_data_tasks/pyproject.toml +++ b/projects/renku_data_tasks/pyproject.toml @@ -44,6 +44,7 @@ packages = [ { include = "renku_data_services/metrics", from = "../../components" }, { include = "renku_data_services/capacity_reservation", from = "../../components" }, { include = "renku_data_services/resource_usage", from = "../../components" }, + { include = "renku_data_services/persisted_logs", from = "../../components" }, ] [tool.poetry.dependencies] diff --git a/pyproject.toml b/pyproject.toml index 807fbf23ff..0d4f83f89b 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -44,6 +44,7 @@ packages = [ { include = "renku_data_services/metrics", from = "components" }, { include = "renku_data_services/notifications", from = "components" }, { include = "renku_data_services/resource_usage", from = "components" }, + { include = "renku_data_services/persisted_logs", from = "components" }, ] [tool.poetry.dependencies] diff --git a/test/components/renku_data_services/persisted_logs/__init__.py b/test/components/renku_data_services/persisted_logs/__init__.py new file mode 100644 index 0000000000..0835fac718 --- /dev/null +++ b/test/components/renku_data_services/persisted_logs/__init__.py @@ -0,0 +1 @@ +"""Tests for persisted logs.""" diff --git a/test/components/renku_data_services/persisted_logs/conftest.py b/test/components/renku_data_services/persisted_logs/conftest.py new file mode 100644 index 0000000000..7e0a256692 --- /dev/null +++ b/test/components/renku_data_services/persisted_logs/conftest.py @@ -0,0 +1,198 @@ +"""Tests for the logs collector.""" + +import pytest + +from renku_data_services.persisted_logs import loki_api + + +@pytest.fixture +def session_logs_response() -> loki_api.LokiQueryRangeResponse: + json_content = """ +{ + "status": "success", + "data": { + "resultType": "streams", + "result": [ + { + "stream": { + "app": "AmaltheaSession", + "container": "git-clone", + "container_runtime": "containerd", + "detected_level": "unknown", + "instance": "renku-ci-ds-1383/j-flora-thie-a8944af936b5-7mjnh:git-clone", + "job": "renku-ci-ds-1383/git-clone", + "namespace": "renku-ci-ds-1383", + "pod": "j-flora-thie-a8944af936b5-7mjnh", + "renku_io_launcher_id": "01KXNAFYMJ42QCEGG6739T28RS", + "renku_io_pod_uid": "c7341da8-5333-47fc-87f4-b364b0ea9c1c", + "renku_io_project_id": "01KXJZF4YH8G2CP9TDNJJPAWNF", + "renku_io_run_id": "01KYVGCNJ3CJ1EQKF343JEEV6T", + "renku_io_safe_username": "d62fb7cb-7893-4149-8917-19e8d882cdd0", + "renku_io_session_type": "non_interactive", + "renku_io_session_uid": "6c5596f3-b27d-4b71-8d37-e672eb66b866", + "renku_io_submission_id": "run-8ej4lg", + "service_name": "AmaltheaSession" + }, + "values": [ + [ + "1785482086378524994", + "2026/07/31 07:14:46 Setting up git proxy to http://localhost:65480\\n" + ], + [ + "1785482086364418022", + "2026/07/31 07:14:46 Dealing with submodules\\n" + ], + [ + "1785482086339950490", + "2026/07/31 07:14:46 Checking out branch main\\n" + ], + [ + "1785482086339935443", + "2026/07/31 07:14:46 Default branch is main\\n" + ], + [ + "1785482085028615483", + "2026/07/31 07:14:45 Cloning repository /home/renku/work/renku-envs from https://gitlab.com/leafty/renku-envs.git\\n" + ], + [ + "1785482085026311874", + "2026/07/31 07:14:45 Setting name Flora Thiebaut in git config\\n" + ], + [ + "1785482085025300000", + "2026/07/31 07:14:45 Setting email flora.thiebaut@sdsc.ethz.ch in git config\\n" + ], + [ + "1785482085020557046", + "2026/07/31 07:14:45 Initializing repo\\n" + ], + [ + "1785482085020548955", + "2026/07/31 07:14:45 Setting up repository.\\n" + ], + [ + "1785482085019553389", + "2026/07/31 07:14:45 Processing https://gitlab.com/leafty/renku-envs.git\\n" + ], + [ + "1785482085018064534", + "2026/07/31 07:14:45 Creating clone path\\n" + ], + [ + "1785482085018035543", + "2026/07/31 07:14:45 Checking if clone path: /home/renku/work/renku-envs exists\\n" + ] + ] + }, + { + "stream": { + "app": "AmaltheaSession", + "container": "amalthea-session", + "container_runtime": "containerd", + "detected_level": "unknown", + "instance": "renku-ci-ds-1383/j-flora-thie-a8944af936b5-7mjnh:amalthea-session", + "job": "renku-ci-ds-1383/amalthea-session", + "namespace": "renku-ci-ds-1383", + "pod": "j-flora-thie-a8944af936b5-7mjnh", + "renku_io_launcher_id": "01KXNAFYMJ42QCEGG6739T28RS", + "renku_io_pod_uid": "c7341da8-5333-47fc-87f4-b364b0ea9c1c", + "renku_io_project_id": "01KXJZF4YH8G2CP9TDNJJPAWNF", + "renku_io_run_id": "01KYVGCNJ3CJ1EQKF343JEEV6T", + "renku_io_safe_username": "d62fb7cb-7893-4149-8917-19e8d882cdd0", + "renku_io_session_type": "non_interactive", + "renku_io_session_uid": "6c5596f3-b27d-4b71-8d37-e672eb66b866", + "renku_io_submission_id": "run-8ej4lg", + "service_name": "AmaltheaSession" + }, + "values": [ + [ + "1785482091170416212", + "10/10\\n" + ], + [ + "1785482091170414253", + "9/10\\n" + ], + [ + "1785482091170412277", + "8/10\\n" + ], + [ + "1785482091170410038", + "7/10\\n" + ], + [ + "1785482091170408130", + "6/10\\n" + ], + [ + "1785482091170406092", + "5/10\\n" + ], + [ + "1785482091170404020", + "4/10\\n" + ], + [ + "1785482091170401840", + "3/10\\n" + ], + [ + "1785482091170397828", + "2/10\\n" + ], + [ + "1785482091170339759", + "1/10\\n" + ] + ] + } + ] + } +} +""" + return loki_api.LokiQueryRangeResponse.model_validate_json(json_content) + + +@pytest.fixture +def build_logs_response() -> loki_api.LokiQueryRangeResponse: + json_content = """ +{ + "status": "success", + "data": { + "resultType": "streams", + "result": [ + { + "stream": { + "app": "ShipwrightBuildRun", + "container": "step-build-and-push", + "container_runtime": "containerd", + "detected_level": "unknown", + "instance": "renku-ci-ds-1383/renku-01kyvgffxtxv4qk0dyjkx0zsa5-ttr4v-pod:step-build-and-push", + "job": "renku-ci-ds-1383/step-build-and-push", + "namespace": "renku-ci-ds-1383", + "pod": "renku-01kyvgffxtxv4qk0dyjkx0zsa5-ttr4v-pod", + "renku_io_buildrun_name": "renku-01kyvgffxtxv4qk0dyjkx0zsa5", + "renku_io_pod_uid": "a533d0d6-c485-4671-b3d6-643ad34bd0b6", + "service_name": "ShipwrightBuildRun" + }, + "values": [ + [ + "1785482346922298906", + " harbor.dev.renku.ch/renku-build/renku-build:renku-01kyvgffxtxv4qk0dyjkx0zsa5\\n" + ], + [ + "1785482346922274287", + "*** Images (sha256:77281bcd4ffcc16bd9942dd280d1a9ef78c64fd97274cd1ec4b0d4a0c4084fef):\\n" + ], + [ + "1785482342006624181", + "Saving harbor.dev.renku.ch/renku-build/renku-build:renku-01kyvgffxtxv4qk0dyjkx0zsa5...\\n" + ] + ] + } + ] + } +} +""" + return loki_api.LokiQueryRangeResponse.model_validate_json(json_content) diff --git a/test/components/renku_data_services/persisted_logs/test_collector.py b/test/components/renku_data_services/persisted_logs/test_collector.py new file mode 100644 index 0000000000..2554f9b9df --- /dev/null +++ b/test/components/renku_data_services/persisted_logs/test_collector.py @@ -0,0 +1,64 @@ +"""Tests for the logs collector.""" + +import pytest +from ulid import ULID + +from renku_data_services.persisted_logs import loki_api, models +from renku_data_services.persisted_logs.collector import LokiLogReader + + +@pytest.mark.asyncio +async def test_process_session_logs(session_logs_response: loki_api.LokiQueryRangeResponse) -> None: + log_stream = LokiLogReader._process_session_logs(session_logs_response) + unsaved_log_lines: list[models.UnsavedSessionLogLine] = [] + async for item in log_stream: + unsaved_log_lines.append(item) + + assert unsaved_log_lines is not None + assert len(unsaved_log_lines) == 22 + + expected_log_line_1 = models.UnsavedSessionLogLine( + id="1785482085020548955::git-clone::j-flora-thie-a8944af936b5-7mjnh", + user_id="d62fb7cb-7893-4149-8917-19e8d882cdd0", + run_id=ULID.from_str("01KYVGCNJ3CJ1EQKF343JEEV6T"), + session_uid="6c5596f3-b27d-4b71-8d37-e672eb66b866", + launcher_id=ULID.from_str("01KXNAFYMJ42QCEGG6739T28RS"), + submission_id="run-8ej4lg", + container="git-clone", + timestamp=1785482085020548955, + log_line="2026/07/31 07:14:45 Setting up repository.\n", + ) + assert expected_log_line_1 in unsaved_log_lines + + expected_log_line_2 = models.UnsavedSessionLogLine( + id="1785482091170410038::amalthea-session::j-flora-thie-a8944af936b5-7mjnh", + user_id="d62fb7cb-7893-4149-8917-19e8d882cdd0", + run_id=ULID.from_str("01KYVGCNJ3CJ1EQKF343JEEV6T"), + session_uid="6c5596f3-b27d-4b71-8d37-e672eb66b866", + launcher_id=ULID.from_str("01KXNAFYMJ42QCEGG6739T28RS"), + submission_id="run-8ej4lg", + container="amalthea-session", + timestamp=1785482091170410038, + log_line="7/10\n", + ) + assert expected_log_line_2 in unsaved_log_lines + + +@pytest.mark.asyncio +async def test_process_build_logs(build_logs_response: loki_api.LokiQueryRangeResponse) -> None: + log_stream = LokiLogReader._process_image_build_logs(build_logs_response) + unsaved_log_lines: list[models.UnsavedBuildLogLine] = [] + async for item in log_stream: + unsaved_log_lines.append(item) + + assert unsaved_log_lines is not None + assert len(unsaved_log_lines) == 3 + + expected_log_line = models.UnsavedBuildLogLine( + id="1785482342006624181::step-build-and-push::renku-01kyvgffxtxv4qk0dyjkx0zsa5-ttr4v-pod", + build_id=ULID.from_str("01KYVGFFXTXV4QK0DYJKX0ZSA5"), + container="step-build-and-push", + timestamp=1785482342006624181, + log_line="Saving harbor.dev.renku.ch/renku-build/renku-build:renku-01kyvgffxtxv4qk0dyjkx0zsa5...\n", + ) + assert expected_log_line in unsaved_log_lines diff --git a/test/components/renku_data_services/persisted_logs/test_db.py b/test/components/renku_data_services/persisted_logs/test_db.py new file mode 100644 index 0000000000..6ee2e1d1bc --- /dev/null +++ b/test/components/renku_data_services/persisted_logs/test_db.py @@ -0,0 +1,281 @@ +"""Tests for the persisted logs database.""" + +from collections.abc import AsyncIterator +from dataclasses import replace + +import pytest +import pytest_asyncio +from sqlalchemy import select +from ulid import ULID + +from renku_data_services import base_models +from renku_data_services.data_api.dependencies import DependencyManager +from renku_data_services.migrations.core import run_migrations_for_app +from renku_data_services.persisted_logs import loki_api, models +from renku_data_services.persisted_logs import orm as schemas +from renku_data_services.persisted_logs.collector import LokiLogReader +from renku_data_services.persisted_logs.db import ( + AmaltheaSessionPersistedLogsWriteRepository, + ImageBuildPersistedLogsWriteRepository, +) +from renku_data_services.project.models import Project, UnsavedProject, Visibility +from renku_data_services.session.models import ( + Build, + LauncherType, + Platform, + SessionLauncher, + UnsavedBuildParameters, + UnsavedSessionLauncher, +) +from renku_data_services.users.models import UserInfo + + +@pytest_asyncio.fixture +async def dependency_manager(app_manager_instance: DependencyManager) -> DependencyManager: + run_migrations_for_app("common") + return app_manager_instance + + +@pytest.fixture +def session_logs_repo() -> AmaltheaSessionPersistedLogsWriteRepository: + return AmaltheaSessionPersistedLogsWriteRepository() + + +@pytest.fixture +def build_logs_repo() -> ImageBuildPersistedLogsWriteRepository: + return ImageBuildPersistedLogsWriteRepository() + + +@pytest_asyncio.fixture +async def regular_user(dependency_manager: DependencyManager) -> base_models.AuthenticatedAPIUser: + api_user = base_models.AuthenticatedAPIUser( + id="jane_doe", email="jane.doe@example.org", access_token="my_access_token" + ) + user_info = await dependency_manager.kc_user_repo.get_or_create_user(requested_by=api_user, id=api_user.id) + assert user_info is not None + return api_user + + +@pytest_asyncio.fixture +async def regular_user_info( + dependency_manager: DependencyManager, regular_user: base_models.AuthenticatedAPIUser +) -> UserInfo: + user_info = await dependency_manager.kc_user_repo.get_user(id=regular_user.id) + assert user_info is not None + return user_info + + +@pytest_asyncio.fixture +async def my_project( + dependency_manager: DependencyManager, + regular_user: base_models.AuthenticatedAPIUser, + regular_user_info: UserInfo, +) -> Project: + project = await dependency_manager.project_repo.insert_project( + user=regular_user, + project=UnsavedProject( + name="My Project", + slug="my-project", + visibility=Visibility.PRIVATE, + created_by=regular_user.id, + namespace=regular_user_info.namespace.path.serialize(), + ), + ) + assert project is not None + return project + + +@pytest_asyncio.fixture +async def my_session_launcher( + dependency_manager: DependencyManager, + regular_user: base_models.AuthenticatedAPIUser, + my_project: Project, +) -> SessionLauncher: + build_parameters = UnsavedBuildParameters( + repository="https://example.org/repo.git", + platforms=[Platform.linux_amd64], + builder_variant="python", + frontend_variant="vscodium", + ) + launcher = await dependency_manager.session_repo.insert_launcher( + user=regular_user, + launcher=UnsavedSessionLauncher( + project_id=my_project.id, + name="My Session", + description=None, + resource_class_id=None, + disk_storage=None, + env_variables=None, + environment=build_parameters, + launcher_type=LauncherType.interactive, + ), + ) + assert launcher is not None + return launcher + + +@pytest_asyncio.fixture +async def my_build( + dependency_manager: DependencyManager, + regular_user: base_models.AuthenticatedAPIUser, + my_session_launcher: SessionLauncher, +) -> Build: + builds = await dependency_manager.session_repo.get_environment_builds( + user=regular_user, environment_id=my_session_launcher.environment.id + ) + assert builds is not None + assert len(builds) == 1 + return builds[0] + + +@pytest.mark.asyncio +async def test_session_latest_log_timestamp_is_none_at_startup( + session_logs_repo: AmaltheaSessionPersistedLogsWriteRepository, dependency_manager: DependencyManager +): + async_session_maker = dependency_manager.config.db.async_session_maker + async with async_session_maker() as session, session.begin(): + ts = await session_logs_repo.get_latest_log_timestamp(session=session) + assert ts is None + + +@pytest.mark.asyncio +async def test_insert_session_logs( + session_logs_response: loki_api.LokiQueryRangeResponse, + session_logs_repo: AmaltheaSessionPersistedLogsWriteRepository, + dependency_manager: DependencyManager, + regular_user: base_models.AuthenticatedAPIUser, + my_session_launcher: SessionLauncher, +): + # Replace the log line metadata for the test + async def _make_logs_stream() -> AsyncIterator[models.UnsavedSessionLogLine]: + source = LokiLogReader._process_session_logs(session_logs_response) + async for item in source: + yield replace(item, user_id=regular_user.id, launcher_id=my_session_launcher.id) + + async_session_maker = dependency_manager.config.db.async_session_maker + async with async_session_maker() as session, session.begin(): + result = await session_logs_repo.insert_session_logs(session=session, logs_stream=_make_logs_stream()) + + expected = models.LogStreamMetadata(log_count=22, last_timestamp=1785482091170416212) + assert result == expected + + # Check the result of get_latest_log_timestamp() + async with async_session_maker() as session, session.begin(): + ts = await session_logs_repo.get_latest_log_timestamp(session=session) + assert ts == expected.last_timestamp + + +@pytest.mark.asyncio +async def test_insert_session_logs_with_db_failures( + session_logs_response: loki_api.LokiQueryRangeResponse, + session_logs_repo: AmaltheaSessionPersistedLogsWriteRepository, + dependency_manager: DependencyManager, + regular_user: base_models.AuthenticatedAPIUser, + my_session_launcher: SessionLauncher, +): + # Replace the log line metadata for the test + async def _make_logs_stream() -> AsyncIterator[models.UnsavedSessionLogLine]: + source = LokiLogReader._process_session_logs(session_logs_response) + alt_run_id = ULID() + idx = 0 + async for item in source: + # # Use correct foreign keys only on half of the log lines + if idx % 2 == 0: + yield replace(item, user_id=regular_user.id, run_id=alt_run_id, launcher_id=my_session_launcher.id) + else: + yield item + idx += 1 + + async_session_maker = dependency_manager.config.db.async_session_maker + async with async_session_maker() as session, session.begin(): + result = await session_logs_repo.insert_session_logs(session=session, logs_stream=_make_logs_stream()) + + expected = models.LogStreamMetadata(log_count=22, last_timestamp=1785482091170416212) + assert result == expected + + async with async_session_maker() as session, session.begin(): + stmt = select(schemas.AmaltheaSessionLogORM) + res = await session.scalars(stmt) + session_logs_orm = res.all() + + assert len(session_logs_orm) == 11 + + # Check the result of get_latest_log_timestamp() + async with async_session_maker() as session, session.begin(): + ts = await session_logs_repo.get_latest_log_timestamp(session=session) + assert ts == expected.last_timestamp + + +@pytest.mark.asyncio +async def test_build_latest_log_timestamp_is_none_at_startup( + build_logs_repo: ImageBuildPersistedLogsWriteRepository, dependency_manager: DependencyManager +): + async_session_maker = dependency_manager.config.db.async_session_maker + async with async_session_maker() as session, session.begin(): + ts = await build_logs_repo.get_latest_log_timestamp(session=session) + assert ts is None + + +@pytest.mark.asyncio +async def test_insert_build_logs( + build_logs_response: loki_api.LokiQueryRangeResponse, + build_logs_repo: ImageBuildPersistedLogsWriteRepository, + dependency_manager: DependencyManager, + my_build: Build, +): + # Replace the log line metadata for the test + async def _make_logs_stream() -> AsyncIterator[models.UnsavedBuildLogLine]: + source = LokiLogReader._process_image_build_logs(build_logs_response) + async for item in source: + yield replace(item, build_id=my_build.id) + + async_session_maker = dependency_manager.config.db.async_session_maker + async with async_session_maker() as session, session.begin(): + result = await build_logs_repo.insert_build_logs(session=session, logs_stream=_make_logs_stream()) + + expected = models.LogStreamMetadata(log_count=3, last_timestamp=1785482346922298906) + assert result == expected + + # Check the result of get_latest_log_timestamp() + async with async_session_maker() as session, session.begin(): + ts = await build_logs_repo.get_latest_log_timestamp(session=session) + assert ts == expected.last_timestamp + + +@pytest.mark.asyncio +async def test_insert_build_logs_with_db_failures( + build_logs_response: loki_api.LokiQueryRangeResponse, + build_logs_repo: ImageBuildPersistedLogsWriteRepository, + dependency_manager: DependencyManager, + my_build: Build, +): + # Replace the log line metadata for the test + async def _make_logs_stream() -> AsyncIterator[models.UnsavedBuildLogLine]: + source = LokiLogReader._process_image_build_logs(build_logs_response) + idx = 0 + async for item in source: + # Use correct foreign keys only on half of the log lines + if idx % 2 == 0: + yield replace(item, build_id=my_build.id) + else: + yield item + idx += 1 + + async_session_maker = dependency_manager.config.db.async_session_maker + async with async_session_maker() as session, session.begin(): + result = await build_logs_repo.insert_build_logs(session=session, logs_stream=_make_logs_stream()) + + expected = models.LogStreamMetadata(log_count=3, last_timestamp=1785482346922298906) + assert result == expected + + async with async_session_maker() as session, session.begin(): + stmt = select(schemas.ImageBuildLogORM) + res = await session.scalars(stmt) + build_logs_orm = res.all() + + assert len(build_logs_orm) == 2 + + # Check the result of get_latest_log_timestamp() + async with async_session_maker() as session, session.begin(): + ts = await build_logs_repo.get_latest_log_timestamp(session=session) + assert ts == expected.last_timestamp diff --git a/test/utils.py b/test/utils.py index 0ad39e4458..dd1c926b3b 100644 --- a/test/utils.py +++ b/test/utils.py @@ -56,6 +56,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, @@ -364,6 +368,12 @@ def from_env( occurrence_repo = OccurrenceRepository(session_maker=config.db.async_session_maker) resource_requests_repo = ResourceRequestsRepo(session_maker=config.db.async_session_maker) resource_usage_service = ResourceUsageService(resource_requests_repo) + 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=config, @@ -410,6 +420,8 @@ def from_env( 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,