From ceeae77e7d683ca939a6b9920a54431c536f19ce Mon Sep 17 00:00:00 2001 From: Travis Herbranson Date: Sun, 24 May 2026 22:46:48 -0400 Subject: [PATCH] postgres migration: schema, models, embeddings, alembic MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Swap the SQLite backing store for petalbrain Postgres + pgvector, modeled on vault-mcp. All second-brain relational tables now live in the `second_brain` schema (owned by the lovebug role); embeddings are written to the shared public.embeddings table. Locked design decisions (per Travis): - DB: existing petalbrain Postgres, second_brain schema, lovebug role. - Connection: containerized homelab-postgres:5432, plain psycopg_pool (min=1/max=10), no PgBouncer. - ORM stays SQLAlchemy; int autoincrement PKs + naive UTC DateTime. - Embeddings: reuse shared public.embeddings keyed by (source_schema='second_brain', source_table='extractions', source_id, model='nomic-embed-text'). Summaries only for this round. - Pipeline: chunk_text → Ollama nomic-embed-text → delete-before-insert upsert, with graceful degradation (no DB / no Ollama → log + skip). - Alembic stands up second-brain's own schema; public.embeddings stays out-of-band. - File-based wiki compiler is unchanged. No SQLite data import — starting clean. This commit is the scaffolding only; `alembic upgrade head` and a smoke test of the embedding path are the next checkpoint. --- alembic.ini | 52 ++++ alembic/env.py | 108 ++++++++ alembic/script.py.mako | 27 ++ ...f6c9d10_v1_second_brain_pipeline_tables.py | 117 ++++++++ config/settings.toml | 33 ++- config/settings.toml.example | 69 +++++ pyproject.toml | 10 + src/second_brain/config.py | 60 ++++- src/second_brain/database.py | 117 ++++++-- src/second_brain/embeddings/__init__.py | 251 ++++++++++++++++++ src/second_brain/main.py | 21 +- src/second_brain/models.py | 36 ++- src/second_brain/scheduler/runner.py | 13 +- 13 files changed, 867 insertions(+), 47 deletions(-) create mode 100644 alembic.ini create mode 100644 alembic/env.py create mode 100644 alembic/script.py.mako create mode 100644 alembic/versions/4a1e2f6c9d10_v1_second_brain_pipeline_tables.py create mode 100644 config/settings.toml.example create mode 100644 src/second_brain/embeddings/__init__.py diff --git a/alembic.ini b/alembic.ini new file mode 100644 index 0000000..bbc8a8e --- /dev/null +++ b/alembic.ini @@ -0,0 +1,52 @@ +# Alembic config for the second_brain schema migrations. +# +# Mirrors vault-mcp/alembic.ini. The sqlalchemy.url placeholder lets the +# `alembic` CLI parse this file without env vars at import time; env.py +# overrides it from SECOND_BRAIN_DATABASE_URL / HERBYLAB_DATABASE_URL at +# runtime. + +[alembic] +script_location = %(here)s/alembic + +prepend_sys_path = . + +path_separator = os + +sqlalchemy.url = postgresql+psycopg://placeholder:placeholder@localhost/placeholder + + +[post_write_hooks] + +[loggers] +keys = root,sqlalchemy,alembic + +[handlers] +keys = console + +[formatters] +keys = generic + +[logger_root] +level = WARNING +handlers = console +qualname = + +[logger_sqlalchemy] +level = WARNING +handlers = +qualname = sqlalchemy.engine + +[logger_alembic] +level = INFO +handlers = +qualname = alembic + +[handler_console] +class = StreamHandler +args = (sys.stderr,) +level = NOTSET +formatter = generic + +[formatter_generic] +format = %(levelname)-5.5s [%(name)s] %(message)s +datefmt = %H:%M:%S diff --git a/alembic/env.py b/alembic/env.py new file mode 100644 index 0000000..fb460c7 --- /dev/null +++ b/alembic/env.py @@ -0,0 +1,108 @@ +"""Alembic environment for second-brain — second_brain schema. + +Mirrors vault-mcp/alembic/env.py. Reads SECOND_BRAIN_DATABASE_URL (with +HERBYLAB_DATABASE_URL as a fallback) so the same migration tree works +across local dev, CI, and production. + +The connection bootstraps the `second_brain` schema and pins alembic's +own version table to that schema; migration scripts themselves set +search_path before issuing unqualified DDL. + +The bootstrap-then-commit pattern below is the fix for the "autobegin +trap": running any further statement before context.begin_transaction() +would autobegin a new transaction, causing alembic's begin_transaction() +to nest as a savepoint that silently rolls back when the connection +closes. Do not regress this. +""" + +from __future__ import annotations + +import os +from logging.config import fileConfig + +from dotenv import load_dotenv +from sqlalchemy import engine_from_config, pool + +from alembic import context + +# Make the project's `src/` importable so `second_brain.models` resolves +# without a `pip install -e .`. +import sys +from pathlib import Path + +_PROJECT_ROOT = Path(__file__).resolve().parent.parent +sys.path.insert(0, str(_PROJECT_ROOT / "src")) + +from second_brain.models import Base, SCHEMA # noqa: E402 + +load_dotenv() + +config = context.config + +if config.config_file_name is not None: + fileConfig(config.config_file_name) + + +def _resolve_db_url() -> str | None: + """Pick the first usable URL out of the two homelab conventions.""" + for var in ("SECOND_BRAIN_DATABASE_URL", "HERBYLAB_DATABASE_URL"): + v = os.environ.get(var) + if v: + return v + return None + + +_db_url = _resolve_db_url() +if _db_url: + config.set_main_option("sqlalchemy.url", _db_url) + +target_metadata = Base.metadata + + +def run_migrations_offline() -> None: + url = config.get_main_option("sqlalchemy.url") + context.configure( + url=url, + target_metadata=target_metadata, + literal_binds=True, + dialect_opts={"paramstyle": "named"}, + version_table_schema=SCHEMA, + include_schemas=True, + ) + + with context.begin_transaction(): + context.run_migrations() + + +def run_migrations_online() -> None: + connectable = engine_from_config( + config.get_section(config.config_ini_section, {}), + prefix="sqlalchemy.", + poolclass=pool.NullPool, + ) + + with connectable.connect() as connection: + # Bootstrap the schema before alembic looks for alembic_version. + # Commit explicitly so the schema persists; otherwise the implicit + # transaction is rolled back when the connection closes. Avoid + # running any further statement here — it would autobegin a new + # transaction and cause alembic's `begin_transaction()` to nest as + # a savepoint, which silently rolls back when the connection closes. + # The migration script sets search_path itself. + connection.exec_driver_sql(f"CREATE SCHEMA IF NOT EXISTS {SCHEMA}") + connection.commit() + context.configure( + connection=connection, + target_metadata=target_metadata, + version_table_schema=SCHEMA, + include_schemas=True, + ) + + with context.begin_transaction(): + context.run_migrations() + + +if context.is_offline_mode(): + run_migrations_offline() +else: + run_migrations_online() diff --git a/alembic/script.py.mako b/alembic/script.py.mako new file mode 100644 index 0000000..891119b --- /dev/null +++ b/alembic/script.py.mako @@ -0,0 +1,27 @@ +"""${message} + +Revision ID: ${up_revision} +Revises: ${down_revision | comma,n} +Create Date: ${create_date} + +""" +from __future__ import annotations + +from collections.abc import Sequence + +from alembic import op +import sqlalchemy as sa +${imports if imports else ""} + +revision: str = ${repr(up_revision)} +down_revision: str | Sequence[str] | None = ${repr(down_revision)} +branch_labels: str | Sequence[str] | None = ${repr(branch_labels)} +depends_on: str | Sequence[str] | None = ${repr(depends_on)} + + +def upgrade() -> None: + ${upgrades if upgrades else "pass"} + + +def downgrade() -> None: + ${downgrades if downgrades else "pass"} diff --git a/alembic/versions/4a1e2f6c9d10_v1_second_brain_pipeline_tables.py b/alembic/versions/4a1e2f6c9d10_v1_second_brain_pipeline_tables.py new file mode 100644 index 0000000..da2d4eb --- /dev/null +++ b/alembic/versions/4a1e2f6c9d10_v1_second_brain_pipeline_tables.py @@ -0,0 +1,117 @@ +"""v1 second_brain: sources, extractions, wiki_pages + +Revision ID: 4a1e2f6c9d10 +Revises: +Create Date: 2026-05-24 00:00:00.000000 + +Emits the v1 second_brain schema: the relational backbone of the +extraction pipeline. Enum types (`source_type`, `source_status`) live +inside the schema so they don't collide with anything else in the +petalbrain cluster. + +env.py creates the second_brain schema and pins alembic_version to it +before this migration runs. Below we re-issue CREATE SCHEMA IF NOT +EXISTS defensively and SET search_path so unqualified names resolve. + +Vector indexing lives in public.embeddings — managed out-of-band, not +by this migration tree. +""" + +from __future__ import annotations + +from collections.abc import Sequence + +from alembic import op + +revision: str = "4a1e2f6c9d10" +down_revision: str | Sequence[str] | None = None +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + + +def upgrade() -> None: + """Create the v1 second_brain schema.""" + op.execute("CREATE SCHEMA IF NOT EXISTS second_brain") + op.execute("SET search_path TO second_brain") + + op.execute( + "CREATE TYPE source_type AS ENUM ('video', 'article')" + ) + op.execute( + """ + CREATE TYPE source_status AS ENUM ( + 'pending', 'pulled', 'transcribed', 'analyzed', + 'accepted', 'published', 'failed' + ) + """ + ) + + op.execute( + """ + CREATE TABLE sources ( + id SERIAL PRIMARY KEY, + url VARCHAR(2000) NOT NULL UNIQUE, + title VARCHAR(500), + source_type source_type NOT NULL DEFAULT 'video', + domain VARCHAR(50) NOT NULL DEFAULT 'development', + focus TEXT, + status source_status NOT NULL DEFAULT 'pending', + media_path VARCHAR(1000), + transcript_path VARCHAR(1000), + transcript_text TEXT, + published_at TIMESTAMP, + duration_seconds INTEGER, + author VARCHAR(200), + ingested_at TIMESTAMP NOT NULL DEFAULT (now() AT TIME ZONE 'UTC'), + updated_at TIMESTAMP, + error_message TEXT + ) + """ + ) + op.execute("CREATE INDEX idx_sources_status ON sources(status)") + op.execute("CREATE INDEX idx_sources_domain ON sources(domain)") + + op.execute( + """ + CREATE TABLE extractions ( + id SERIAL PRIMARY KEY, + source_id INTEGER NOT NULL UNIQUE + REFERENCES sources(id) ON DELETE CASCADE, + summary TEXT, + key_points JSON, + entities JSON, + claims JSON, + open_questions JSON, + contradictions JSON, + raw_response TEXT, + input_tokens INTEGER, + output_tokens INTEGER, + created_at TIMESTAMP NOT NULL DEFAULT (now() AT TIME ZONE 'UTC') + ) + """ + ) + + op.execute( + """ + CREATE TABLE wiki_pages ( + id SERIAL PRIMARY KEY, + vault_path VARCHAR(500) NOT NULL UNIQUE, + domain VARCHAR(50) NOT NULL, + title VARCHAR(300) NOT NULL, + git_sha_before VARCHAR(40), + git_sha_after VARCHAR(40), + created_at TIMESTAMP NOT NULL DEFAULT (now() AT TIME ZONE 'UTC'), + updated_at TIMESTAMP + ) + """ + ) + + +def downgrade() -> None: + """Drop the v1 second_brain tables and enums (schema is left intact).""" + op.execute("SET search_path TO second_brain") + op.execute("DROP TABLE IF EXISTS wiki_pages") + op.execute("DROP TABLE IF EXISTS extractions") + op.execute("DROP TABLE IF EXISTS sources") + op.execute("DROP TYPE IF EXISTS source_status") + op.execute("DROP TYPE IF EXISTS source_type") diff --git a/config/settings.toml b/config/settings.toml index 84fcc23..08d4657 100644 --- a/config/settings.toml +++ b/config/settings.toml @@ -2,10 +2,9 @@ # Copy this file and edit paths to match your setup. # Override the config file path with SECOND_BRAIN_CONFIG env var. -# Path to the SQLite database -db_path = "~/.local/share/second-brain/second_brain.db" - -# Path to your Obsidian vault (git-managed directory) +# Path to your Obsidian vault (git-managed directory). The wiki compiler +# writes markdown into this directory and commits inside it. We do NOT +# dual-write into petalbrain.wiki — the file-based vault stays canonical. vault_path = "~/Documents/second-brain-vault" # Directory where downloaded media files are stored @@ -21,6 +20,21 @@ whisper_model = "small" # Active knowledge domains domains = ["development", "content", "business", "homelab"] +[database] +# Postgres + pgvector lives on the petalbrain instance. +# +# Preferred: leave `url` empty and export SECOND_BRAIN_DATABASE_URL (or the +# shared HERBYLAB_DATABASE_URL) in the environment so the password never +# lands in this checked-in file. +# +# Containerised second-brain runs should target homelab-postgres on the +# homelab docker network: +# postgresql+psycopg://lovebug:@homelab-postgres:5432/petalbrain +# +# Host-side dev runs can point at the published port instead: +# postgresql+psycopg://lovebug:@127.0.0.1:5433/petalbrain +url = "" + [extractor] # "cli" → runs `claude -p` under your Max OAuth (no per-token billing). # "api" → uses Anthropic SDK, requires ANTHROPIC_API_KEY. @@ -29,6 +43,17 @@ model = "claude-sonnet-4-6" cli_binary = "claude" timeout_sec = 600 +[embeddings] +# Best-effort: extraction summaries are chunked (512 tokens / 64 overlap) +# and embedded via Ollama's nomic-embed-text (768-d) into the shared +# public.embeddings table. Failures (no DB / no Ollama) log a warning and +# the pipeline keeps moving — the relational + file vault stay canonical. +enabled = true +# Containerized default; for host-side runs set OLLAMA_URL=http://127.0.0.1:11434 +ollama_url = "http://ollama:11434" +model = "nomic-embed-text" +embed_field = "summary" + [scheduler] # Time window for overnight processing (24h format, may cross midnight) window_start = "22:00" diff --git a/config/settings.toml.example b/config/settings.toml.example new file mode 100644 index 0000000..08d4657 --- /dev/null +++ b/config/settings.toml.example @@ -0,0 +1,69 @@ +# second-brain runtime configuration +# Copy this file and edit paths to match your setup. +# Override the config file path with SECOND_BRAIN_CONFIG env var. + +# Path to your Obsidian vault (git-managed directory). The wiki compiler +# writes markdown into this directory and commits inside it. We do NOT +# dual-write into petalbrain.wiki — the file-based vault stays canonical. +vault_path = "~/Documents/second-brain-vault" + +# Directory where downloaded media files are stored +media_dir = "~/.local/share/second-brain/media" + +# Directory where Whisper SRT transcripts are stored +subtitles_dir = "~/.local/share/second-brain/subtitles" + +# Whisper model size: tiny | base | small | medium | large +# small is a good balance of speed and accuracy +whisper_model = "small" + +# Active knowledge domains +domains = ["development", "content", "business", "homelab"] + +[database] +# Postgres + pgvector lives on the petalbrain instance. +# +# Preferred: leave `url` empty and export SECOND_BRAIN_DATABASE_URL (or the +# shared HERBYLAB_DATABASE_URL) in the environment so the password never +# lands in this checked-in file. +# +# Containerised second-brain runs should target homelab-postgres on the +# homelab docker network: +# postgresql+psycopg://lovebug:@homelab-postgres:5432/petalbrain +# +# Host-side dev runs can point at the published port instead: +# postgresql+psycopg://lovebug:@127.0.0.1:5433/petalbrain +url = "" + +[extractor] +# "cli" → runs `claude -p` under your Max OAuth (no per-token billing). +# "api" → uses Anthropic SDK, requires ANTHROPIC_API_KEY. +backend = "cli" +model = "claude-sonnet-4-6" +cli_binary = "claude" +timeout_sec = 600 + +[embeddings] +# Best-effort: extraction summaries are chunked (512 tokens / 64 overlap) +# and embedded via Ollama's nomic-embed-text (768-d) into the shared +# public.embeddings table. Failures (no DB / no Ollama) log a warning and +# the pipeline keeps moving — the relational + file vault stay canonical. +enabled = true +# Containerized default; for host-side runs set OLLAMA_URL=http://127.0.0.1:11434 +ollama_url = "http://ollama:11434" +model = "nomic-embed-text" +embed_field = "summary" + +[scheduler] +# Time window for overnight processing (24h format, may cross midnight) +window_start = "22:00" +window_end = "06:00" + +# CLI backend: cap calls per rolling hour (Max sub has 5h-window message caps). +max_calls_per_hour = 30 + +# Minimum seconds between extraction calls (rate-limit smoother) +min_gap_seconds = 120 + +# Legacy API-backend budget (only enforced when extractor.backend = "api"). +max_tokens_per_hour = 100_000 diff --git a/pyproject.toml b/pyproject.toml index 83c0d8f..c927da2 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,10 +4,15 @@ version = "0.1.0" description = "Personal knowledge system: video/article extraction pipeline + LLM-maintained wiki" requires-python = ">=3.12" dependencies = [ + "alembic>=1.13.0", "click>=8.1.0", + "embedding-chunking", "fastapi>=0.115.0", + "httpx>=0.28.0", "jinja2>=3.1.0", "openai-whisper", + "psycopg[binary]>=3.2", + "psycopg-pool>=3.2", "pydantic>=2.0.0", "python-dotenv>=1.0.0", "sqlalchemy>=2.0.0", @@ -23,6 +28,11 @@ dependencies = [ # Max OAuth subscription and has no Python-side Anthropic dependency. api = ["anthropic>=0.40.0"] +[tool.uv.sources] +# Shared chunker that vault-mcp / ob1-enricher / ob1-reflector also pin. +# Same source so the 512/64 chunking stays in lockstep across services. +embedding-chunking = { git = "https://gitea.plantbasedsoutherner.com/petal-power/embedding-chunking.git", tag = "v0.1.0" } + [project.scripts] second-brain = "second_brain.main:cli" diff --git a/src/second_brain/config.py b/src/second_brain/config.py index 5d8a738..242f3c0 100644 --- a/src/second_brain/config.py +++ b/src/second_brain/config.py @@ -3,6 +3,11 @@ Configuration loading for second-brain. Settings are read from config/settings.toml (or a path set via SECOND_BRAIN_CONFIG). All values have sensible defaults so the system works out of the box. + +The relational store is Postgres (petalbrain DB, `second_brain` schema). Connection +details are resolved by `second_brain.database` — see that module for the URL +resolution order. The vault is still a file-based Obsidian directory; this +project does NOT dual-write to petalbrain.wiki. """ from __future__ import annotations @@ -21,12 +26,17 @@ except ModuleNotFoundError: # --------------------------------------------------------------------------- _DEFAULTS: dict[str, Any] = { - "db_path": "~/.local/share/second-brain/second_brain.db", "vault_path": "~/Documents/second-brain-vault", "media_dir": "~/.local/share/second-brain/media", "subtitles_dir": "~/.local/share/second-brain/subtitles", "whisper_model": "small", "domains": ["development", "content", "business", "homelab"], + "database": { + # Empty by default — the runtime expects SECOND_BRAIN_DATABASE_URL + # (or HERBYLAB_DATABASE_URL) in the environment. Fill this only if + # you prefer keeping the URL inside settings.toml. + "url": "", + }, "extractor": { # "cli" runs `claude -p` under the Max OAuth subscription (default). # "api" uses the Anthropic SDK and requires ANTHROPIC_API_KEY. @@ -35,6 +45,19 @@ _DEFAULTS: dict[str, Any] = { "cli_binary": "claude", "timeout_sec": 600, }, + "embeddings": { + # Best-effort embedding of extraction summaries into public.embeddings. + # If disabled, or if Ollama / the DB is unreachable, the pipeline + # logs a warning and keeps going — file/relational data stays the + # source of truth. + "enabled": True, + # Default targets the `ollama` service on the homelab docker network + # (same convention vault-mcp uses). For host-side runs, set + # OLLAMA_URL=http://127.0.0.1:11434 in the env instead. + "ollama_url": "http://ollama:11434", + "model": "nomic-embed-text", + "embed_field": "summary", + }, "scheduler": { "window_start": "22:00", "window_end": "06:00", @@ -49,6 +72,13 @@ _DEFAULTS: dict[str, Any] = { DOMAINS = ["development", "content", "business", "homelab"] +def _merge(default: dict, override: dict) -> dict: + """Shallow merge — override wins per top-level key.""" + out = dict(default) + out.update(override or {}) + return out + + class Config: """Holds resolved configuration values.""" @@ -57,10 +87,6 @@ class Config: # --- paths --- - @property - def db_path(self) -> Path: - return Path(self._raw.get("db_path", _DEFAULTS["db_path"])).expanduser() - @property def vault_path(self) -> Path: return Path(self._raw.get("vault_path", _DEFAULTS["vault_path"])).expanduser() @@ -89,17 +115,23 @@ class Config: def anthropic_api_key(self) -> str | None: return os.getenv("ANTHROPIC_API_KEY") - # --- extractor --- + # --- structured sections --- + + @property + def database(self) -> dict[str, Any]: + return _merge(_DEFAULTS["database"], self._raw.get("database", {})) @property def extractor(self) -> dict[str, Any]: - return {**_DEFAULTS["extractor"], **self._raw.get("extractor", {})} + return _merge(_DEFAULTS["extractor"], self._raw.get("extractor", {})) - # --- scheduler --- + @property + def embeddings(self) -> dict[str, Any]: + return _merge(_DEFAULTS["embeddings"], self._raw.get("embeddings", {})) @property def scheduler(self) -> dict[str, Any]: - return {**_DEFAULTS["scheduler"], **self._raw.get("scheduler", {})} + return _merge(_DEFAULTS["scheduler"], self._raw.get("scheduler", {})) # --- prompts dir --- @@ -112,7 +144,7 @@ class Config: def ensure_dirs(self) -> None: """Create runtime directories if they don't exist.""" - for d in (self.db_path.parent, self.media_dir, self.subtitles_dir): + for d in (self.media_dir, self.subtitles_dir): d.mkdir(parents=True, exist_ok=True) @@ -150,4 +182,10 @@ def load_config(path: Path | None = None) -> Config: return _config_instance -__all__ = ["Config", "DOMAINS", "load_config"] +def reset_config_singleton() -> None: + """Test hook — clear the cached Config so a re-read picks up new settings.""" + global _config_instance + _config_instance = None + + +__all__ = ["Config", "DOMAINS", "load_config", "reset_config_singleton"] diff --git a/src/second_brain/database.py b/src/second_brain/database.py index 7c91f55..71050d0 100644 --- a/src/second_brain/database.py +++ b/src/second_brain/database.py @@ -1,43 +1,105 @@ """ Database engine and session management for second-brain. -Ported and adapted from xtract/database.py. -Uses a singleton pattern: call get_database() everywhere. +Postgres + pgvector backend on the petalbrain instance. + +Resolution order for the connection URL: +1. constructor arg passed to `get_database()` +2. env var `SECOND_BRAIN_DATABASE_URL` +3. env var `HERBYLAB_DATABASE_URL` (shared homelab convention) +4. `[database].url` in settings.toml + +The URL is expected to use SQLAlchemy's `postgresql+psycopg://` driver +(psycopg v3). Plain `postgresql://` URLs are accepted and rewritten to use +psycopg; this matches the vault-mcp convention where one URL feeds both +Alembic (`+psycopg`) and raw psycopg (no prefix). + +Connection mode: containerized (homelab-postgres:5432 on the homelab docker +network). For local-host dev runs, point at 127.0.0.1:5433 instead. """ from __future__ import annotations +import os from contextlib import contextmanager -from pathlib import Path from typing import Generator -from sqlalchemy import create_engine +from sqlalchemy import create_engine, text +from sqlalchemy.engine import Engine from sqlalchemy.orm import Session, sessionmaker -from second_brain.models import Base +from second_brain.models import SCHEMA, Base + + +def _normalize_url(url: str) -> str: + """Ensure the URL uses SQLAlchemy's psycopg-v3 driver.""" + if url.startswith("postgresql+psycopg://"): + return url + if url.startswith("postgresql://"): + return url.replace("postgresql://", "postgresql+psycopg://", 1) + return url + + +def _resolve_url(explicit: str | None) -> str: + if explicit: + return _normalize_url(explicit) + for var in ("SECOND_BRAIN_DATABASE_URL", "HERBYLAB_DATABASE_URL"): + v = os.environ.get(var) + if v: + return _normalize_url(v) + + # Fall back to settings.toml [database].url + try: + from second_brain.config import load_config + + cfg = load_config() + url = cfg.database.get("url") + if url: + return _normalize_url(url) + except Exception: + pass + + raise RuntimeError( + "No Postgres connection URL found. Set SECOND_BRAIN_DATABASE_URL " + "(or HERBYLAB_DATABASE_URL), or fill `[database].url` in " + "config/settings.toml." + ) class Database: - """Manages the SQLAlchemy engine and session factory.""" + """Manages the SQLAlchemy engine and session factory for second-brain.""" - def __init__(self, db_path: Path) -> None: - self.db_path = db_path - self.db_path.parent.mkdir(parents=True, exist_ok=True) + def __init__(self, url: str | None = None) -> None: + self.url = _resolve_url(url) - self.engine = create_engine( - f"sqlite:///{self.db_path}", + # Pool sizing matches vault-mcp (min=1/max=10). No PgBouncer; a normal + # QueuePool is fine for the CLI + scheduler + web workers we run. + self.engine: Engine = create_engine( + self.url, echo=False, - connect_args={"check_same_thread": False}, + pool_size=int(os.environ.get("SECOND_BRAIN_DB_POOL_MIN", "1")), + max_overflow=int(os.environ.get("SECOND_BRAIN_DB_POOL_MAX", "9")), + pool_pre_ping=True, + future=True, ) self.SessionLocal = sessionmaker( autocommit=False, autoflush=False, bind=self.engine, + future=True, ) def init_db(self) -> None: - """Create all tables (idempotent).""" + """Create the schema and tables if missing. + + Schema creation is idempotent and only relevant for dev/test setups + where Alembic hasn't been applied yet. In production, prefer + `alembic upgrade head` — this method just defensively ensures the + schema exists before `create_all` issues unqualified DDL. + """ + with self.engine.begin() as conn: + conn.execute(text(f"CREATE SCHEMA IF NOT EXISTS {SCHEMA}")) Base.metadata.create_all(bind=self.engine) def get_session(self) -> Session: @@ -57,6 +119,9 @@ class Database: finally: sess.close() + def close(self) -> None: + self.engine.dispose() + # --------------------------------------------------------------------------- # Singleton @@ -65,17 +130,25 @@ class Database: _db_instance: Database | None = None -def get_database(db_path: Path | None = None) -> Database: - """Return the singleton Database, creating and initialising it on first call.""" +def get_database(url: str | None = None) -> Database: + """Return the singleton Database, creating it lazily on first call. + + Schema/tables are *not* auto-created here — run `alembic upgrade head` + once during setup. Tests may explicitly call `Database(url).init_db()` + against a scratch database. + """ global _db_instance if _db_instance is None: - if db_path is None: - from second_brain.config import load_config - - db_path = load_config().db_path - _db_instance = Database(db_path) - _db_instance.init_db() + _db_instance = Database(url) return _db_instance -__all__ = ["Database", "get_database"] +def reset_database_singleton() -> None: + """Test hook — drop the cached singleton so a new URL can take effect.""" + global _db_instance + if _db_instance is not None: + _db_instance.close() + _db_instance = None + + +__all__ = ["Database", "get_database", "reset_database_singleton"] diff --git a/src/second_brain/embeddings/__init__.py b/src/second_brain/embeddings/__init__.py new file mode 100644 index 0000000..61b43b9 --- /dev/null +++ b/src/second_brain/embeddings/__init__.py @@ -0,0 +1,251 @@ +"""Embed second-brain extractions into the shared public.embeddings table. + +This is a near-verbatim copy of the vault-mcp recipe (vault_mcp.core.embeddings), +adapted for our `source_schema='second_brain' / source_table='extractions'` +key tuple and our embed payload (the LLM-generated summary). + +Behaviour: +- Reuses the shared `embedding_chunking` library for 512-token chunking with + 64-token overlap. +- Calls Ollama's /api/embeddings with `num_ctx=8192` to embed nomic-embed-text + vectors (768-d). +- Writes vector literals as `%s::vector` to avoid taking a hard dependency on + pgvector's Python adapter. +- Delete-before-insert on the (schema, table, source_id, embedding_model) + tuple so a re-embed of an extraction that previously had N chunks but now + has M str | None: + """Resolve a raw psycopg DSN (no SQLAlchemy `+psycopg` prefix).""" + for var in ("SECOND_BRAIN_DATABASE_URL", "HERBYLAB_DATABASE_URL"): + v = os.environ.get(var) + if v: + return v.replace("postgresql+psycopg://", "postgresql://", 1) + + # Fallback: pull from settings.toml so a host-side `process` works even + # when the operator forgot to export the env var. + try: + from second_brain.config import load_config + + url = load_config().database.get("url") + if url: + return url.replace("postgresql+psycopg://", "postgresql://", 1) + except Exception: + pass + return None + + +def _get_pool() -> ConnectionPool | None: + """Lazy-init a small connection pool. Returns None if no DB URL is configured.""" + global _pool + if _pool is not None: + return _pool + dsn = _database_url() + if not dsn: + return None + _pool = ConnectionPool(dsn, min_size=1, max_size=4, open=True, timeout=10) + return _pool + + +def close_pool() -> None: + """Close the embedding-side pool. Safe to call multiple times.""" + global _pool + if _pool is not None: + _pool.close() + _pool = None + + +# --------------------------------------------------------------------------- +# Embedding model client +# --------------------------------------------------------------------------- + + +def _ollama_url() -> str: + """Env var takes precedence over settings.toml so containerized runs can + point at a different Ollama instance without rebuilding the image.""" + env = os.environ.get("OLLAMA_URL") + if env: + return env + try: + from second_brain.config import load_config + + return load_config().embeddings.get("ollama_url", "http://ollama:11434") + except Exception: + return "http://ollama:11434" + + +def _embedding_model() -> str: + env = os.environ.get("EMBEDDING_MODEL") + if env: + return env + try: + from second_brain.config import load_config + + return load_config().embeddings.get("model", "nomic-embed-text") + except Exception: + return "nomic-embed-text" + + +def embed_text(text: str) -> list[float]: + """Call Ollama's embeddings endpoint. Returns the model-native vector.""" + resp = httpx.post( + f"{_ollama_url()}/api/embeddings", + json={ + "model": _embedding_model(), + "prompt": text, + # nomic-embed-text supports 8192-token context. Ollama defaults + # the runner to 2048, which 500s on long summaries. + "options": {"num_ctx": 8192}, + }, + timeout=120.0, + ) + resp.raise_for_status() + data = resp.json() + vec = data.get("embedding") + if not vec or not isinstance(vec, list): + raise RuntimeError( + f"Ollama returned no embedding for model {_embedding_model()}" + ) + return vec + + +def _vector_literal(vec: list[float]) -> str: + """pgvector accepts text-literal vectors of form '[v1,v2,...]'.""" + return "[" + ",".join(repr(float(v)) for v in vec) + "]" + + +# --------------------------------------------------------------------------- +# Public API +# --------------------------------------------------------------------------- + + +def embed_extraction(extraction_id: int, summary: str | None) -> int | None: + """Embed an extraction summary into public.embeddings. + + Args: + extraction_id: PK of the row in second_brain.extractions. + summary: the LLM-generated summary string; we embed this and + nothing else for the current round (key_points / claims + are not embedded yet). + + Returns the number of chunk rows written, or None on skip/failure. + Best-effort: failures are logged, not raised. + """ + if not summary or not summary.strip(): + logger.debug("embed_extraction: skip — extraction %s has no summary", extraction_id) + return None + + pool = _get_pool() + if pool is None: + logger.info( + "embed_extraction: skip — no SECOND_BRAIN_DATABASE_URL configured" + ) + return None + + model = _embedding_model() + + try: + chunks: list[Chunk] = chunk_text(summary) + chunk_embeddings: list[tuple[Chunk, list[float]]] = [ + (c, embed_text(c.text)) for c in chunks + ] + except Exception as exc: + logger.warning( + "embed_extraction: chunk/embed failed for extraction %s: %s", + extraction_id, + exc, + ) + return None + + try: + with pool.connection() as conn: + n = _upsert_embeddings( + conn, extraction_id, chunk_embeddings, model + ) + conn.commit() + logger.info( + "embed_extraction: extraction=%s chunks=%d model=%s", + extraction_id, + n, + model, + ) + return n + except Exception as exc: + logger.warning( + "embed_extraction: DB write failed for extraction %s: %s", + extraction_id, + exc, + ) + return None + + +def _upsert_embeddings( + conn: Any, + extraction_id: int, + chunk_embeddings: list[tuple[Chunk, list[float]]], + model: str, +) -> int: + """Replace all public.embeddings rows for this (extraction, model) pair.""" + with conn.cursor() as cur: + cur.execute( + """ + DELETE FROM public.embeddings + WHERE source_schema = %s + AND source_table = %s + AND source_id = %s + AND embedding_model = %s + """, + (SOURCE_SCHEMA, SOURCE_TABLE, extraction_id, model), + ) + for chunk, vec in chunk_embeddings: + cur.execute( + """ + INSERT INTO public.embeddings + (source_schema, source_table, source_id, chunk_index, + chunk_text, chunk_token_count, embedding, embedding_model) + VALUES (%s, %s, %s, %s, %s, %s, %s::vector, %s) + """, + ( + SOURCE_SCHEMA, + SOURCE_TABLE, + extraction_id, + chunk.index, + chunk.text, + chunk.token_count, + _vector_literal(vec), + model, + ), + ) + return len(chunk_embeddings) + + +__all__ = ["embed_extraction", "embed_text", "close_pool"] diff --git a/src/second_brain/main.py b/src/second_brain/main.py index 065e98d..7a5dab2 100644 --- a/src/second_brain/main.py +++ b/src/second_brain/main.py @@ -54,7 +54,7 @@ def add(url: str, domain: str, focus: Optional[str], title: Optional[str]) -> No """Queue a source URL for processing.""" config = load_config() config.ensure_dirs() - db = get_database(config.db_path) + db = get_database() source_type = SourceType.ARTICLE if _is_article(url) else SourceType.VIDEO @@ -115,7 +115,7 @@ def process( config = load_config() config.ensure_dirs() - db = get_database(config.db_path) + db = get_database() yt_adapter = YouTubeAdapter(config) art_adapter = ArticleAdapter(config) @@ -199,19 +199,32 @@ def process( existing_ext.raw_response = getattr(result, "_raw_response", None) existing_ext.input_tokens = getattr(result, "_input_tokens", None) existing_ext.output_tokens = getattr(result, "_output_tokens", None) + ext_row = existing_ext else: - extraction = Extraction( + ext_row = Extraction( source_id=source.id, raw_response=getattr(result, "_raw_response", None), input_tokens=getattr(result, "_input_tokens", None), output_tokens=getattr(result, "_output_tokens", None), **result.to_db_dict(), ) - sess.add(extraction) + sess.add(ext_row) source.status = SourceStatus.ANALYZED source.updated_at = utcnow() + # Flush so the new extraction gets its PK before we kick + # off embedding (we key public.embeddings rows by it). + sess.flush() + extraction_id = ext_row.id + extraction_summary = ext_row.summary click.echo(" ✓ analyzed") + + if config.embeddings.get("enabled", True): + from second_brain.embeddings import embed_extraction + + n = embed_extraction(extraction_id, extraction_summary) + if n is not None: + click.echo(f" ✓ embedded ({n} chunk(s))") except Exception as exc: source.error_message = str(exc) source.status = SourceStatus.FAILED diff --git a/src/second_brain/models.py b/src/second_brain/models.py index dd4d6f7..0454350 100644 --- a/src/second_brain/models.py +++ b/src/second_brain/models.py @@ -1,6 +1,10 @@ """ SQLAlchemy ORM models for second-brain. +All tables live in the Postgres `second_brain` schema (in the petalbrain DB). +PK style is int autoincrement (SERIAL in Postgres); timestamps are naive UTC +(`timestamp without time zone`) per the locked design decisions in CLAUDE.md. + Tables: - Source — a queued/processed URL (video or article) - Extraction — structured LLM output for a source @@ -24,14 +28,19 @@ from sqlalchemy import ( ForeignKey, Integer, JSON, + MetaData, String, Text, ) from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column, relationship +SCHEMA = "second_brain" + class Base(DeclarativeBase): - """Shared declarative base.""" + """Shared declarative base. All tables live in the `second_brain` schema.""" + + metadata = MetaData(schema=SCHEMA) # --------------------------------------------------------------------------- @@ -54,6 +63,22 @@ class SourceStatus(enum.Enum): FAILED = "failed" +# Native Postgres enum types, scoped to the second_brain schema so they +# don't collide with anything else in the cluster. +_SourceTypeEnum = Enum( + SourceType, + name="source_type", + schema=SCHEMA, + values_callable=lambda et: [e.value for e in et], +) +_SourceStatusEnum = Enum( + SourceStatus, + name="source_status", + schema=SCHEMA, + values_callable=lambda et: [e.value for e in et], +) + + # --------------------------------------------------------------------------- # Models # --------------------------------------------------------------------------- @@ -69,7 +94,7 @@ class Source(Base): url: Mapped[str] = mapped_column(String(2000), unique=True, nullable=False) title: Mapped[Optional[str]] = mapped_column(String(500), nullable=True) source_type: Mapped[SourceType] = mapped_column( - Enum(SourceType), nullable=False, default=SourceType.VIDEO + _SourceTypeEnum, nullable=False, default=SourceType.VIDEO ) # Domain and focus @@ -78,7 +103,7 @@ class Source(Base): # Pipeline status status: Mapped[SourceStatus] = mapped_column( - Enum(SourceStatus), nullable=False, default=SourceStatus.PENDING + _SourceStatusEnum, nullable=False, default=SourceStatus.PENDING ) # File paths (set after pull/transcribe steps) @@ -115,7 +140,9 @@ class Extraction(Base): id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) source_id: Mapped[int] = mapped_column( - ForeignKey("sources.id", ondelete="CASCADE"), nullable=False, unique=True + ForeignKey(f"{SCHEMA}.sources.id", ondelete="CASCADE"), + nullable=False, + unique=True, ) # Core fields (mirrors the extraction schema contract) @@ -167,6 +194,7 @@ class WikiPage(Base): __all__ = [ "Base", + "SCHEMA", "Source", "SourceStatus", "SourceType", diff --git a/src/second_brain/scheduler/runner.py b/src/second_brain/scheduler/runner.py index ee63aca..6541692 100644 --- a/src/second_brain/scheduler/runner.py +++ b/src/second_brain/scheduler/runner.py @@ -127,18 +127,27 @@ class Scheduler: existing.raw_response = getattr(result, "_raw_response", None) existing.input_tokens = getattr(result, "_input_tokens", None) existing.output_tokens = getattr(result, "_output_tokens", None) + ext_row = existing else: - sess.add(Extraction( + ext_row = Extraction( source_id=source.id, raw_response=getattr(result, "_raw_response", None), input_tokens=getattr(result, "_input_tokens", None), output_tokens=getattr(result, "_output_tokens", None), **fields, - )) + ) + sess.add(ext_row) source.status = SourceStatus.ANALYZED source.updated_at = utcnow() + # Flush so the embedding step can key on the new PK. + sess.flush() + if self.config.embeddings.get("enabled", True): + from second_brain.embeddings import embed_extraction + + embed_extraction(ext_row.id, ext_row.summary) + calls_this_hour += 1 used = (getattr(result, "_input_tokens", 0) or 0) + ( getattr(result, "_output_tokens", 0) or 0