All verification complete — TODO.md was already cleared in the roadmap commit; the two extra unit-test diffs are necessary fake-signature adaptations for the new keywords. Everything is green, no fixes were needed. ## Phase 89 — final verification pass: ALL GREEN **Verified (all 6 task files present in `complete/`):** - `git_sources.ignore_paths` JSONB column + migration 0013; `alembic downgrade -1 && upgrade head` round-trips (head `0013`) - Importer: `normalize_ignore_path`/`is_ignored`/`_ignore_for_root`, `ignore` in walk + progress pre-walk, `ignore_by_root` in `import_sources` - API: GET/POST carry list; admin-only `PATCH` (replace, 404/422 fixed details, anonymous 403) - Pipelines wired: `_run_sync`, `_run_upload` re-upload, `scripts/import_docs.py` - Sources-page box: dialog, §7.4 save lifecycle, `N ignored` tag, a11y; env rows get no box **Test/lint results:** - `uv run pytest` → 1808 passed - `uv run pytest --cov=app --cov-report=term-missing` → TOTAL **99%** (>90%) - `uv run pytest tests/e2e/test_source_ignore_paths.py -v --no-cov` → 6 passed (isolated, DB up) - Regressions in isolation: `test_git_sources_admin` 6, `test_archive_upload_sources` 5, `test_sync_button` 3, `test_smoke` 3 — all passed - `uv run ruff check . && uv run pyright` → clean (0 errors) **Completion criteria:** box→PATCH 200→count+GET round-trip ✅ · sync excludes `ignore/` (no docs/chunks/embeddings/summaries) + prunes newly-ignored (pruned==2) ✅ · no-mid-path rule E2E ✅ · PATCH 404/422/replace/clear/403 ✅ · full gate green ✅ · commit + phase move left to harness per rules. **Deviations:** none blocking — E2E pins `files == 4` (overview's "5" was an off-by-one vs its own 6-file tree, documented in-test); `tests/unit/test_importer.py` + `test_sync_button.py` test-double fakes extended for the new keywords (needed for the suite to stay green). **Next pending phase:** none — `todo/` holds only this phase.
795 lines
35 KiB
Python
795 lines
35 KiB
Python
"""Admin-managed sources API (phase 35, task 02; local kind, phase 38).
|
||
|
||
Admin-only CRUD under ``/api/git-sources`` (phase 16 pattern, A10
|
||
extended — the public API surface stays stateless and the signed cookie
|
||
remains the only session state, same as ``/api/steering`` and
|
||
``/api/sync``): the ``git_sources`` table holds the sources the Sync
|
||
button (phase 32) and ``import_docs`` (phase 28) import — ``kind='git'``
|
||
rows carry the repo URL to clone/pull, ``kind='local'`` rows (phase 38)
|
||
carry an existing directory on the server to walk directly. DB rows win
|
||
over ``BOR_GIT_SOURCES``, which is a git-only fallback while the table is
|
||
empty (the phase's locked decision — ``from_env`` tells the UI which list
|
||
it is looking at, so the page can show the env note only while the
|
||
fallback is active).
|
||
|
||
Routes: ``GET`` (DB rows oldest-first, or the env list with
|
||
``from_env: true`` while the table is empty; rows carry ``kind`` +
|
||
``path``, git rows — and env rows — report ``path: null``; every row
|
||
reports ``ignore_paths`` — DB rows their stored normalized list, env
|
||
rows ``[]``, phase 89), ``POST``
|
||
(201, validated create; ``kind`` selects the validation: git → exactly
|
||
the phase-35 URL contract, local → an existing absolute directory, else
|
||
422 naming the path; both kinds accept ``ignore_paths`` — optional,
|
||
absent → ``[]``, stored normalized, phase 89),
|
||
``PATCH /{source_id}`` (phase 89, A5 — replace one source's ignore
|
||
list: 404 unknown id, the required body list is normalized + A4-
|
||
validated with fixed-detail 422s and REPLACES the row's list wholesale
|
||
— an empty list clears all; 200 → the ``GitSourceOut`` shape),
|
||
``POST /upload`` (phase 49, backgrounded in phase
|
||
64 task 03 — admin archive upload: the ``.tar``/``.tar.gz``/``.tgz``/
|
||
``.zip`` name/format gate + the 1 MiB-chunk receive with the
|
||
``upload_max_mb`` cap run **inline** and answered 202 the moment the
|
||
archive is safely on disk; unpack → swap → row upsert → model check →
|
||
single-source scan → change-gated overview then run in a **background
|
||
task** — see :func:`upload_archive` and :func:`_run_upload`),
|
||
``GET /upload/status`` (the phase-32 ``SyncStatus``-shaped in-memory
|
||
state of that run — incl. the phase-64 ``current_file`` /
|
||
``files_done`` / ``files_total`` progress fields; navigating away from
|
||
the page mid-scan no longer aborts anything), ``DELETE /{source_id}``
|
||
(204 — total removal, phase 69: row + the source's documents (chunks +
|
||
embeddings) committed first, then the app-managed on-disk dir). The
|
||
whole router sits behind :func:`app.core.auth.require_admin` —
|
||
anonymous callers get 403 on every route.
|
||
|
||
No credential-echo path: git URLs may embed ``user:pass@`` (phase 32's
|
||
masking discipline), so every git 409/422 detail is a fixed generic
|
||
string that never repeats the submitted URL. Local paths are not
|
||
secrets — the local 422/409 details name the (expanded) path so the
|
||
owner sees exactly which directory failed.
|
||
|
||
Scope boundary (phase locked decisions): the CRUD routes do NOT
|
||
clone or import anything — the existing Sync button performs that, and
|
||
its ``prune=True`` stays for **upstream file churn** (files deleted
|
||
upstream or dropped from a local dir), not for row removal: ``DELETE``
|
||
is a total removal in itself (phase 69) — the row + the source's
|
||
documents (chunks + embeddings via the ``all, delete-orphan`` cascade)
|
||
commit atomically **first** (the RAG is always consistent with the
|
||
registry), then the app-managed on-disk dir (git checkout / unpacked
|
||
upload folder) is deleted; an ``OSError`` there is logged, never fatal.
|
||
Foreign local dirs (the owner's own) are never touched on disk, a
|
||
sibling row sharing the source name keeps the shared documents +
|
||
files (only the row goes), and a pruned KB bumps ``sources_version``
|
||
exactly once (the phase-53 saved-chat invalidation) with a
|
||
best-effort overview refresh. The upload route is the other exception
|
||
(phase 64, task 03): after the 202 receive
|
||
answer, its background task unpacks the archive, swaps it in, upserts
|
||
the row, probes the models, scans the single source
|
||
(``import_sources`` with ``prune=True`` + the change-gated overview
|
||
refresh), and lands the sync-style counts (the ``UploadOut`` fields)
|
||
in the status ``detail``.
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import logging
|
||
import re
|
||
import shutil
|
||
import time
|
||
import uuid
|
||
from dataclasses import dataclass, field
|
||
from datetime import UTC, datetime
|
||
from pathlib import Path
|
||
from typing import Any, Literal, cast
|
||
|
||
from fastapi import APIRouter, Depends, File, HTTPException, Response, UploadFile
|
||
from sqlalchemy import select
|
||
from sqlalchemy.exc import IntegrityError
|
||
from sqlalchemy.orm import Session
|
||
|
||
from app.api.sync import _sanitize_error
|
||
from app.config import get_settings
|
||
from app.core.auth import require_admin
|
||
from app.db import SessionLocal, get_db
|
||
from app.models import Document, GitSource
|
||
from app.rag.archive_upload import (
|
||
ARCHIVE_SUFFIXES,
|
||
ArchiveUploadError,
|
||
archive_source_name,
|
||
swap_in,
|
||
unpack_archive,
|
||
)
|
||
from app.rag.importer import import_sources, normalize_ignore_path
|
||
from app.rag.llm import LLMClient, check_models
|
||
from app.rag.overview import regenerate_overview
|
||
from app.rag.source_removal import (
|
||
has_sibling,
|
||
managed_dir_for,
|
||
remove_managed_dir,
|
||
resolve_source_name,
|
||
)
|
||
from app.rag.sources_meta import bump_sources_version
|
||
from app.schemas import (
|
||
GitSourceIgnoreIn,
|
||
GitSourceIn,
|
||
GitSourceList,
|
||
GitSourceOut,
|
||
GitSourceRow,
|
||
UploadAccepted,
|
||
)
|
||
|
||
logger = logging.getLogger("app.api.git_sources")
|
||
|
||
router = APIRouter(
|
||
prefix="/git-sources",
|
||
tags=["git-sources"],
|
||
dependencies=[Depends(require_admin)], # phase 16 pattern: admin-only surface
|
||
)
|
||
|
||
#: One upload at a time (phase 49, backgrounded in phase 64 task 03):
|
||
#: the flag is checked and set with **no await in between, BEFORE the
|
||
#: streaming receive** — the handler now awaits (the 1 MiB-chunk stream)
|
||
#: long before the background task exists, so a task-done check alone
|
||
#: would let a concurrent POST slip through the receive window and start
|
||
#: a second run. A plain bool, not an ``asyncio.Lock``: it is checked
|
||
#: and set with no await in between (a single app loop can never enter
|
||
#: twice), and it stays correct across requests that run on separate
|
||
#: event loops (the TestClient convention). Held until
|
||
#: :func:`_run_upload`'s ``finally`` (end of the background run) — or
|
||
#: cleared on the inline exception path where the task was never
|
||
#: created.
|
||
_upload_in_progress = False
|
||
|
||
#: Streaming read size while counting compressed upload bytes (1 MiB
|
||
#: chunks — the task-02 cap check granularity).
|
||
_STREAM_CHUNK = 1 << 20
|
||
|
||
|
||
@dataclass
|
||
class UploadStatus:
|
||
"""In-memory state of the (at most one) in-flight upload run.
|
||
|
||
Mirrors :class:`app.api.sync.SyncStatus` (the phase-32 pattern,
|
||
phase 64 task 03): ``state`` is the same four-state machine
|
||
(``idle`` / ``running`` / ``success`` / ``failed``); terminal states
|
||
carry the run's ``detail`` (success — the ``UploadOut`` fields) or
|
||
``error`` (failure — sanitized) so the UI can render the last result
|
||
after a page reload (the re-attach behavior, task 05).
|
||
|
||
Phase 64 (task 03) progress fields: ``current_file`` is the
|
||
``source/relative/path`` the scan is processing right now (null
|
||
outside the import phase — unpack/swap/row/model-check first — and
|
||
in terminal states); ``files_done`` / ``files_total`` carry the
|
||
hook's done/total position and survive a terminal state (the run's
|
||
last position is useful context next to the error).
|
||
"""
|
||
|
||
state: Literal["idle", "running", "success", "failed"] = "idle"
|
||
started_at: datetime | None = None
|
||
finished_at: datetime | None = None
|
||
current_file: str | None = None
|
||
files_done: int = 0
|
||
files_total: int = 0
|
||
detail: dict[str, Any] = field(default_factory=dict)
|
||
error: str | None = None
|
||
|
||
|
||
_upload_status = UploadStatus()
|
||
|
||
#: Accepted git URL shapes — the trimmed URL must *start* with one of them.
|
||
#: Covers the phase-28 real URLs (HTTPS + ``git@`` SSH); scp-style
|
||
#: ``host:repo`` is deliberately rejected (422). ASSUMPTION (task 02): the
|
||
#: accepted shapes are exactly these four prefixes.
|
||
URL_RE = re.compile(r"^(https?://|ssh://|git@)")
|
||
|
||
#: Phase 89, A4 — the per-source ignore-list limits, enforced in
|
||
#: :func:`_validate_ignore_paths` (shared by POST and PATCH): at most
|
||
#: 200 entries, each ≤500 chars after normalization. The 422 details
|
||
#: are fixed strings that never echo the input (the router's
|
||
#: credential-safety discipline, applied for consistency).
|
||
MAX_IGNORE_PATHS = 200
|
||
MAX_IGNORE_PATH_LENGTH = 500
|
||
|
||
|
||
@router.get("", response_model=GitSourceList)
|
||
def list_git_sources(
|
||
db: Session = Depends(get_db), # noqa: B008
|
||
) -> GitSourceList:
|
||
"""The effective source list (git + local rows, phase 38).
|
||
|
||
DB rows ordered by ``(added_at, id)`` (oldest first, id tie-break for
|
||
same-timestamp inserts) with ``from_env: false`` — each row carries
|
||
its ``kind``, its ``ignore_paths`` (phase 89 — the stored,
|
||
normalized list; ``or []`` guards a row that predated the column),
|
||
and, for local rows, the stored ``path`` (git rows and env rows
|
||
report ``path: null``); while the table is empty, the
|
||
``BOR_GIT_SOURCES`` env URLs as git rows (the env fallback is
|
||
git-only) with null ``id``/``added_at``, ``ignore_paths: []`` (no
|
||
DB row to store a list on), and ``from_env: true``.
|
||
"""
|
||
rows = db.scalars(
|
||
select(GitSource).order_by(GitSource.added_at.asc(), GitSource.id.asc())
|
||
).all()
|
||
if rows:
|
||
return GitSourceList(
|
||
sources=[
|
||
# ``ck_git_sources_kind`` (migration 0007) guarantees the
|
||
# value is 'git' or 'local' — the cast documents that.
|
||
GitSourceRow(
|
||
id=row.id,
|
||
kind=cast(Literal["git", "local"], row.kind),
|
||
url=row.url,
|
||
path=row.path,
|
||
added_at=row.added_at,
|
||
ignore_paths=row.ignore_paths or [],
|
||
)
|
||
for row in rows
|
||
],
|
||
from_env=False,
|
||
)
|
||
return GitSourceList(
|
||
sources=[
|
||
GitSourceRow(
|
||
id=None, kind="git", url=url, path=None, added_at=None, ignore_paths=[]
|
||
)
|
||
for url in get_settings().git_source_list
|
||
],
|
||
from_env=True,
|
||
)
|
||
|
||
|
||
@router.post("", response_model=GitSourceOut, status_code=201)
|
||
def create_git_source(
|
||
payload: GitSourceIn,
|
||
db: Session = Depends(get_db), # noqa: B008
|
||
) -> GitSourceOut:
|
||
"""Store one source (fields already trimmed by the schema).
|
||
|
||
``kind="git"`` (default) — exactly the phase-35 contract: 422 when
|
||
the URL shape is not one of the accepted prefixes (generic detail —
|
||
the input is never echoed), 409 when the trimmed URL is already
|
||
stored (the unique index is the backstop against a concurrent insert
|
||
the pre-check missed), 201 + the created row otherwise.
|
||
|
||
``kind="local"`` — ``path`` must expand (``~``) to an absolute,
|
||
existing directory on the server: 422 naming the path otherwise
|
||
(fail loud at add-time — the owner sees it immediately), 409 when
|
||
the path is already stored (detail names the path), 201 + the stored
|
||
row otherwise (``url`` holds the expanded path — the table's
|
||
NOT-NULL location column).
|
||
|
||
Wrong field combinations (git without url, local without path, both
|
||
kinds' fields) are 422 with fixed, input-free details.
|
||
|
||
``ignore_paths`` (phase 89) — optional, both kinds: the raw box
|
||
lines are normalized + A4-validated (``_validate_ignore_paths`` —
|
||
the fixed-detail 422s) and the normalized list is what is stored and
|
||
reported.
|
||
"""
|
||
row = _create_git_row(payload, db) if payload.kind == "git" else _create_local_row(payload, db)
|
||
return GitSourceOut(
|
||
id=row.id, url=row.url, added_at=row.added_at, ignore_paths=row.ignore_paths
|
||
)
|
||
|
||
|
||
def _validate_ignore_paths(raw: list[str] | None) -> list[str]:
|
||
"""Normalize + enforce the A4 limits (phase 89); the fixed 422
|
||
details never echo the input (the credential-safety discipline,
|
||
applied for consistency).
|
||
|
||
Shared by POST and PATCH. The empty check runs FIRST: a
|
||
whitespace-only entry must 422, not be silently dropped (the UI
|
||
drops blank lines client-side; the API stays defensive).
|
||
"""
|
||
entries = [normalize_ignore_path(e) for e in (raw or [])]
|
||
if any(not e for e in entries):
|
||
raise HTTPException(status_code=422, detail="ignore paths must be non-empty")
|
||
if len(entries) > MAX_IGNORE_PATHS:
|
||
raise HTTPException(status_code=422, detail="a source has at most 200 ignore paths")
|
||
if any(len(e) > MAX_IGNORE_PATH_LENGTH for e in entries):
|
||
raise HTTPException(status_code=422, detail="an ignore path exceeds 500 characters")
|
||
return entries
|
||
|
||
|
||
def _commit_new(row: GitSource, duplicate_detail: str, db: Session) -> GitSource:
|
||
"""Insert ``row``; the unique index is the backstop — a concurrent
|
||
insert the pre-check missed still yields the generic 409, never a
|
||
500 (phase-35 convention, now shared by both kinds)."""
|
||
db.add(row)
|
||
try:
|
||
db.commit()
|
||
except IntegrityError:
|
||
db.rollback()
|
||
raise HTTPException(status_code=409, detail=duplicate_detail) from None
|
||
db.refresh(row)
|
||
return row
|
||
|
||
|
||
def _create_git_row(payload: GitSourceIn, db: Session) -> GitSource:
|
||
"""``kind=git`` — the phase-35 URL contract, unchanged (A10: no
|
||
credential echo, so every detail is a fixed string)."""
|
||
if payload.path is not None:
|
||
raise HTTPException(status_code=422, detail="a git source takes a url, not a path")
|
||
if payload.url is None:
|
||
raise HTTPException(status_code=422, detail="a git source requires a url")
|
||
url = payload.url
|
||
if not URL_RE.match(url):
|
||
raise HTTPException(
|
||
status_code=422, detail="not a valid git URL (expected https://, ssh:// or git@…)"
|
||
)
|
||
if db.scalar(select(GitSource).where(GitSource.url == url)) is not None:
|
||
raise HTTPException(status_code=409, detail="a git source with this URL already exists")
|
||
return _commit_new(
|
||
GitSource(
|
||
url=url,
|
||
kind="git",
|
||
ignore_paths=_validate_ignore_paths(payload.ignore_paths),
|
||
),
|
||
"a git source with this URL already exists",
|
||
db,
|
||
)
|
||
|
||
|
||
def _create_local_row(payload: GitSourceIn, db: Session) -> GitSource:
|
||
"""``kind=local`` — fail-loud add-time validation (phase 38):
|
||
trimmed → ``expanduser()`` → absolute + existing directory, else 422
|
||
naming the path (not a secret, unlike a git URL)."""
|
||
if payload.url is not None:
|
||
raise HTTPException(status_code=422, detail="a local source takes a path, not a url")
|
||
if payload.path is None:
|
||
raise HTTPException(status_code=422, detail="a local source requires a path")
|
||
expanded = Path(payload.path).expanduser()
|
||
if not expanded.is_absolute() or not expanded.is_dir():
|
||
raise HTTPException(
|
||
status_code=422, detail=f"local source path is not a directory: {expanded}"
|
||
)
|
||
path = str(expanded)
|
||
if db.scalar(select(GitSource).where(GitSource.path == path)) is not None:
|
||
raise HTTPException(
|
||
status_code=409, detail=f"a local source with this path already exists: {path}"
|
||
)
|
||
# ``url`` is the table's NOT-NULL location column (phase 38: local
|
||
# rows carry the expanded path there too — git URL shapes and absolute
|
||
# paths cannot collide).
|
||
return _commit_new(
|
||
GitSource(
|
||
url=path,
|
||
kind="local",
|
||
path=path,
|
||
ignore_paths=_validate_ignore_paths(payload.ignore_paths),
|
||
),
|
||
f"a local source with this path already exists: {path}",
|
||
db,
|
||
)
|
||
|
||
|
||
@router.patch("/{source_id}", response_model=GitSourceOut)
|
||
def patch_git_source(
|
||
source_id: uuid.UUID,
|
||
payload: GitSourceIgnoreIn,
|
||
db: Session = Depends(get_db), # noqa: B008
|
||
) -> GitSourceOut:
|
||
"""Replace one source's ignore list (phase 89, A5).
|
||
|
||
404 unknown id. The body list (required) is normalized +
|
||
A4-validated (fixed 422 details) and REPLACES the row's list
|
||
wholesale — an empty list clears all. Returns the updated
|
||
row's public shape (id, url, added_at, ignore_paths).
|
||
"""
|
||
row = db.get(GitSource, source_id)
|
||
if row is None:
|
||
raise HTTPException(status_code=404, detail="git source not found")
|
||
row.ignore_paths = _validate_ignore_paths(payload.ignore_paths)
|
||
db.commit()
|
||
db.refresh(row)
|
||
return GitSourceOut(
|
||
id=row.id, url=row.url, added_at=row.added_at, ignore_paths=row.ignore_paths
|
||
)
|
||
|
||
|
||
@router.post("/upload", response_model=UploadAccepted, status_code=202)
|
||
async def upload_archive(
|
||
file: UploadFile = File(...), # noqa: B008
|
||
) -> UploadAccepted:
|
||
"""Receive a source archive; scan it in the background (phase 49,
|
||
backgrounded in phase 64 task 03 — owner-locked A1/A2).
|
||
|
||
The **inline (request) work is exactly three gates** — steps 1–3 —
|
||
everything else runs in a background task behind
|
||
``GET /upload/status`` (the phase-32 ``SyncStatus`` pattern), so
|
||
navigating away mid-scan no longer aborts anything:
|
||
|
||
1. name/format gate — only ``.tar``/``.tar.gz``/``.tgz``/``.zip``
|
||
(422 naming the accepted set) and a safe source name
|
||
(``archive_source_name`` — its message is the 422 detail);
|
||
2. one at a time — 409 ``an upload is already in progress`` while
|
||
the flag is held (checked and set with no await in between,
|
||
BEFORE the receive — see ``_upload_in_progress``);
|
||
3. stream the upload in 1 MiB chunks into a dotfile temp with the
|
||
``upload_max_mb`` cap — 413 naming the cap, temp deleted.
|
||
|
||
Then the archive is **safely on disk** — 202 + ``UploadAccepted``
|
||
(the "successfully uploaded" moment the UI toasts on, A2) and
|
||
:func:`_run_upload` runs the rest on the app's event loop:
|
||
|
||
4. unpack to a temp sibling (traversal/symlink/device/corrupt/
|
||
over-cap → ``failed`` with the task-01 user-safe message, temps
|
||
deleted); a zero-entry archive is ``failed`` ``the archive
|
||
contains no files`` — an archive with only non-A9 files is a
|
||
VALID replacement (the scan indexes nothing, prune removes the
|
||
source's docs);
|
||
5. atomic swap-in — a same-name re-upload replaces the previous
|
||
folder in place; a failure leaves the previous folder/row/KB
|
||
untouched;
|
||
6. upsert the row by ``path`` (``kind='local'``; an existing row is
|
||
left as-is — ``added_at`` preserved — and the unique index is
|
||
the backstop: a concurrent insert lands ``failed`` with
|
||
``a local source with this path already exists: <path>``); the
|
||
row's saved ``ignore_paths`` are captured for the scan (phase
|
||
89: a re-upload of an existing source honors the list the owner
|
||
already saved);
|
||
7. fail-fast ``check_models`` — ``ModelUnavailableError`` →
|
||
``failed`` with the sanitized message (the phase-49 503 becomes
|
||
a status state, A5); the folder/row are already committed, so
|
||
the next sync/re-upload retries idempotently;
|
||
8. ``import_sources([folder], llm, prune=True, progress=<hook>,
|
||
ignore_by_root={folder: row's list})`` (phase 89) + the
|
||
change-gated ``regenerate_overview`` — the hook feeds the status
|
||
``current_file`` / ``files_done`` / ``files_total``;
|
||
9. one INFO log line (PLAN §9 / AGENTS.md rule 10 — ``total_ms`` is
|
||
the background run's duration);
|
||
10. ``success`` — ``detail`` = the ``UploadOut`` fields.
|
||
"""
|
||
# 1. Name/format gate — the accepted formats first (the 422 names
|
||
# them), then the task-01 safe-name derivation. A BARE suffix
|
||
# ("tar.gz") is an accepted format with no usable stem — it
|
||
# passes here and gets task-01's "no usable source name" 422.
|
||
# No upload dir is created for a rejected name.
|
||
filename = file.filename or ""
|
||
lowered = filename.lower()
|
||
if not any(
|
||
lowered.endswith(suffix) or lowered == suffix.lstrip(".")
|
||
for suffix in ARCHIVE_SUFFIXES
|
||
):
|
||
raise HTTPException(
|
||
status_code=422,
|
||
detail="only .tar, .tar.gz, .tgz or .zip archives are accepted",
|
||
)
|
||
try:
|
||
name = archive_source_name(filename)
|
||
except ArchiveUploadError as e:
|
||
raise HTTPException(status_code=422, detail=str(e)) from None
|
||
|
||
# 2. One at a time — the flag is checked and set with no await
|
||
# between, BEFORE the streaming receive: the background task
|
||
# does not exist yet, so the flag (not a task-done check) is the
|
||
# gate (see ``_upload_in_progress``).
|
||
global _upload_in_progress
|
||
if _upload_in_progress:
|
||
raise HTTPException(status_code=409, detail="an upload is already in progress")
|
||
_upload_in_progress = True
|
||
|
||
settings = get_settings()
|
||
upload_root = Path(settings.upload_dir).expanduser()
|
||
upload_root.mkdir(parents=True, exist_ok=True)
|
||
max_bytes = settings.upload_max_mb * 1024 * 1024
|
||
temp_upload = upload_root / f".{name}.{uuid.uuid4().hex}.upload"
|
||
temp_unpack = upload_root / f".{name}.{uuid.uuid4().hex}.unpack"
|
||
try:
|
||
# 3. Stream with the compressed-size cap — dotfile temps are
|
||
# hidden from the upload dir's listing.
|
||
total = 0
|
||
with open(temp_upload, "wb") as out:
|
||
while chunk := await file.read(_STREAM_CHUNK):
|
||
total += len(chunk)
|
||
if total > max_bytes:
|
||
raise HTTPException(
|
||
status_code=413,
|
||
detail=f"the upload exceeds the {settings.upload_max_mb} MiB limit",
|
||
)
|
||
out.write(chunk)
|
||
# The archive is safely on disk — 202 is the "successfully
|
||
# uploaded" moment (A2). Steps 4–10 run in the background:
|
||
asyncio.create_task(
|
||
_run_upload(name, filename, total, upload_root, temp_upload, temp_unpack)
|
||
)
|
||
except BaseException:
|
||
# The background task was never created (cap 413, a broken
|
||
# pipe, cancellation, or create_task itself): release the flag
|
||
# so the next upload is not refused, and make sure no temp
|
||
# survives the failed receive.
|
||
temp_upload.unlink(missing_ok=True)
|
||
_upload_in_progress = False
|
||
raise
|
||
return UploadAccepted(name=name)
|
||
|
||
|
||
@router.get("/upload/status")
|
||
def upload_status() -> dict[str, Any]:
|
||
"""Current upload state (the UI polls this — the phase-32
|
||
``GET /api/sync/status`` contract, identical key set).
|
||
|
||
``started_at`` / ``finished_at`` are ISO-8601 strings or null.
|
||
``current_file`` (phase 64) is the ``source/relative/path`` the
|
||
scan is processing right now — null during the unpack/swap/row/
|
||
model phases and in terminal states; ``files_done`` / ``files_total``
|
||
carry the hook's position (0/0 idle). The router dependency makes
|
||
it admin-only like every other route here.
|
||
"""
|
||
return {
|
||
"state": _upload_status.state,
|
||
"started_at": (
|
||
_upload_status.started_at.isoformat() if _upload_status.started_at else None
|
||
),
|
||
"finished_at": (
|
||
_upload_status.finished_at.isoformat() if _upload_status.finished_at else None
|
||
),
|
||
"detail": _upload_status.detail,
|
||
"error": _upload_status.error,
|
||
"current_file": _upload_status.current_file,
|
||
"files_done": _upload_status.files_done,
|
||
"files_total": _upload_status.files_total,
|
||
}
|
||
|
||
|
||
async def _run_upload(
|
||
name: str,
|
||
filename: str,
|
||
total_bytes: int,
|
||
upload_root: Path,
|
||
temp_upload: Path,
|
||
temp_unpack: Path,
|
||
) -> None:
|
||
"""The post-202 upload pipeline, one in-process background task
|
||
(the phase-32 ``_run_sync`` shape — A1).
|
||
|
||
Every failure mode (unpack, zero entries, swap, row, models,
|
||
import, anything else) lands in the ``failed`` state with a
|
||
sanitized ``error`` string — a background task must die in state,
|
||
never as an unobserved exception (A5: post-202 failures are status
|
||
states, never HTTP errors). ``CancelledError`` is deliberately *not*
|
||
caught: app shutdown cancels the task, and swallowing that would
|
||
mask a real stop. The ``finally`` cleans both temps (defensive —
|
||
each step already cleans its own) and clears ``_upload_in_progress``.
|
||
"""
|
||
global _upload_in_progress
|
||
started = time.monotonic()
|
||
_upload_status.state = "running"
|
||
_upload_status.started_at = datetime.now(UTC)
|
||
_upload_status.finished_at = None
|
||
_upload_status.current_file = None
|
||
_upload_status.files_done = 0
|
||
_upload_status.files_total = 0
|
||
_upload_status.detail = {}
|
||
_upload_status.error = None
|
||
try:
|
||
settings = get_settings()
|
||
max_bytes = settings.upload_max_mb * 1024 * 1024
|
||
# Step 4 — unpack to a temp sibling; the compressed bytes are
|
||
# no longer needed once unpacked (phase 49 locked decision:
|
||
# only the unpacked content is kept).
|
||
unpack_archive(temp_upload, temp_unpack, max_bytes)
|
||
temp_upload.unlink(missing_ok=True)
|
||
if not any(temp_unpack.iterdir()):
|
||
# Zero entries = a user error. (Only non-A9 files is NOT an
|
||
# error — it still has entries and is a valid replacement.)
|
||
raise ArchiveUploadError("the archive contains no files")
|
||
# Step 5 — swap in — a same-name re-upload replaces the
|
||
# previous folder atomically; a failure leaves it, the row,
|
||
# and the KB untouched (the ``failed`` state carries the
|
||
# user-safe message).
|
||
final_dir = upload_root / name
|
||
swap_in(temp_unpack, final_dir)
|
||
# Step 6 — upsert the row by path in a SHORT-LIVED session
|
||
# (open/close around it — the ``effective_sources`` /
|
||
# ``bump_sources_version`` pattern in ``app.api.sync``): the
|
||
# background task has no request session to leak locks from
|
||
# (the old inline ``db.close()`` discipline, now structural).
|
||
# No duplicates: an existing row is left exactly as it is
|
||
# (``added_at`` preserved); the unique index is the backstop
|
||
# for a concurrent insert the pre-check missed.
|
||
path = str(final_dir)
|
||
db = SessionLocal()
|
||
try:
|
||
row = db.scalar(select(GitSource).where(GitSource.path == path))
|
||
if row is None:
|
||
row = GitSource(url=path, kind="local", path=path)
|
||
db.add(row)
|
||
try:
|
||
db.commit()
|
||
except IntegrityError:
|
||
db.rollback()
|
||
raise ValueError(
|
||
f"a local source with this path already exists: {path}"
|
||
) from None
|
||
# Phase 89: the row's saved ignore list, copied to plain
|
||
# values while the row is still usable in this session — a
|
||
# re-upload of an existing source honors the list the owner
|
||
# already saved; a fresh row has no list yet.
|
||
ignore_paths = list(row.ignore_paths or [])
|
||
finally:
|
||
db.close()
|
||
# Step 7 — fail-fast models (phase 41): ``ModelUnavailableError``
|
||
# lands in the ``failed`` state sanitized (the phase-49 503
|
||
# becomes a status state, A5). Nothing is rolled back — the
|
||
# folder/row are committed and the next sync/re-upload retries
|
||
# idempotently.
|
||
llm = LLMClient()
|
||
await check_models(llm)
|
||
# Step 8 — scan — single source, prune (dropped files leave
|
||
# the KB), with the phase-64 progress hook feeding the status,
|
||
# then the change-gated overview refresh (phases 31/32). The
|
||
# closure captures the module ``_upload_status`` exactly like
|
||
# the state assignments above.
|
||
def _hook(source: str, rel: str, done: int, total: int) -> None:
|
||
_upload_status.current_file = f"{source}/{rel}"
|
||
_upload_status.files_done = done
|
||
_upload_status.files_total = total
|
||
|
||
summary = await import_sources(
|
||
[final_dir], llm, prune=True, progress=_hook,
|
||
ignore_by_root={str(final_dir): ignore_paths},
|
||
)
|
||
overview = False
|
||
if summary.added + summary.updated > 0:
|
||
overview = await regenerate_overview(llm)
|
||
# Step 9 — per-upload log line (PLAN §9 / AGENTS.md rule 10)
|
||
# — moved with the scan: ``total_ms`` is the background run's
|
||
# duration.
|
||
logger.info(
|
||
"upload: name=%s file=%s bytes_in=%d files=%d added=%d updated=%d "
|
||
"unchanged=%d pruned=%d errors=%d overview=%s total_ms=%d",
|
||
name,
|
||
filename,
|
||
total_bytes,
|
||
summary.files,
|
||
summary.added,
|
||
summary.updated,
|
||
summary.unchanged,
|
||
summary.pruned,
|
||
summary.errors,
|
||
overview,
|
||
round((time.monotonic() - started) * 1000),
|
||
)
|
||
# Step 10 — success: the ``UploadOut`` fields ride in the
|
||
# status ``detail`` (the UI renders the same result line from
|
||
# the status that the sync button renders from its own).
|
||
_upload_status.state = "success"
|
||
_upload_status.finished_at = datetime.now(UTC)
|
||
_upload_status.current_file = None # phase 64: keep the final counts
|
||
_upload_status.detail = {
|
||
"source": name,
|
||
"files": summary.files,
|
||
"added": summary.added,
|
||
"updated": summary.updated,
|
||
"unchanged": summary.unchanged,
|
||
"pruned": summary.pruned,
|
||
"errors": summary.errors,
|
||
"chunks": summary.chunks,
|
||
"overview": overview,
|
||
}
|
||
except Exception as e: # noqa: BLE001 — a background task dies in state, see above
|
||
logger.exception("upload: failed")
|
||
_upload_status.state = "failed"
|
||
_upload_status.finished_at = datetime.now(UTC)
|
||
_upload_status.error = _sanitize_error(str(e))
|
||
_upload_status.current_file = None # phase 64: keep the final counts
|
||
finally:
|
||
_upload_in_progress = False
|
||
# No temp may survive any failure path (defensive — each step
|
||
# already cleans its own; on success both are already gone).
|
||
temp_upload.unlink(missing_ok=True)
|
||
shutil.rmtree(temp_unpack, ignore_errors=True)
|
||
|
||
|
||
@router.delete("/{source_id}", status_code=204)
|
||
async def delete_git_source(
|
||
source_id: uuid.UUID,
|
||
db: Session = Depends(get_db), # noqa: B008
|
||
) -> Response:
|
||
"""Remove a stored source — **totally** (phase 69); 404 unknown id.
|
||
|
||
Total removal, one action: the row, every indexed document of the
|
||
source (chunks + embeddings via the ``all, delete-orphan`` cascade
|
||
— ``app.models``), and, for app-managed sources, the on-disk dir
|
||
(the git checkout or the unpacked upload folder). Foreign local
|
||
directories (the owner's own) are never touched on disk — their row
|
||
+ index entries still go.
|
||
|
||
Locked order (phase 69):
|
||
|
||
1. **Sibling guard** — another stored row resolves to the same
|
||
source name (e.g. ``…/r`` and ``…/r.git``): it still owns the
|
||
shared documents + files, so only this row is deleted (logged
|
||
loudly); no index/disk work.
|
||
2. **DB first** — in the request transaction: the source's
|
||
documents are deleted, then the row, one ``commit``. A DB
|
||
failure propagates as 500 **before any disk work** — the 204
|
||
contract below never sees a half-removal.
|
||
3. **Disk second** — the app-managed dir is removed after the
|
||
commit; an ``OSError`` is logged, never fatal (a leftover dir
|
||
is inert and self-heals on re-add; the reverse order is
|
||
forbidden — a disk failure must never leave a row pointing at
|
||
deleted files). Foreign local dir → skipped.
|
||
4. **When documents were pruned** — the best-effort overview
|
||
refresh (an LLM failure logs, never fails the 204; the next
|
||
added/updated change refreshes it, as today) and then exactly
|
||
one ``sources_version`` bump (phase 53 — the bump lands even if
|
||
the overview failed, mirroring ``_run_sync``'s order). No docs
|
||
pruned → no overview, no bump.
|
||
|
||
The route contract is unchanged: 204 with no body; the per-
|
||
operation INFO line (PLAN §9) carries the counts —
|
||
``files_removed`` is ``yes|no|skipped`` (``skipped`` for the
|
||
sibling guard and for foreign local dirs).
|
||
"""
|
||
started = time.monotonic()
|
||
row = db.get(GitSource, source_id)
|
||
if row is None:
|
||
raise HTTPException(status_code=404, detail="git source not found")
|
||
|
||
name = resolve_source_name(row)
|
||
settings = get_settings()
|
||
managed = managed_dir_for(row, Path(settings.sources_dir), Path(settings.upload_dir))
|
||
|
||
if has_sibling(db, row):
|
||
# Sibling guard: the shared documents + files still belong to
|
||
# the sibling row — only this row goes (phase 69 locked
|
||
# decision), so no index prune and no disk work.
|
||
logger.warning(
|
||
"source removal: row %s (url=%s) shares source name %s with "
|
||
"another stored row — deleting only the row; the shared "
|
||
"documents and files stay",
|
||
row.id,
|
||
row.url,
|
||
name,
|
||
)
|
||
db.delete(row)
|
||
db.commit()
|
||
docs_pruned, files_removed, overview = 0, "skipped", False
|
||
else:
|
||
# DB first: prune every document of the source (the cascade
|
||
# drops all chunks incl. embeddings) and the row itself in one
|
||
# commit — the RAG is always consistent with the registry.
|
||
docs = db.scalars(select(Document).where(Document.source == name)).all()
|
||
for doc in docs:
|
||
db.delete(doc)
|
||
db.delete(row)
|
||
db.commit()
|
||
docs_pruned = len(docs)
|
||
# Disk second: only the app-managed dir (git checkout / upload
|
||
# folder) — ``managed is None`` is the foreign local dir (the
|
||
# owner's own), never touched on disk.
|
||
if managed is None:
|
||
files_removed = "skipped"
|
||
else:
|
||
files_removed = "yes" if remove_managed_dir(managed) else "no"
|
||
# KB changed (documents pruned) → best-effort overview refresh,
|
||
# then exactly one sources_version bump — the bump in its own
|
||
# short-lived session (the ``_run_sync`` pattern) so it lands
|
||
# even if the best-effort overview failed.
|
||
overview = False
|
||
if docs_pruned > 0:
|
||
try:
|
||
overview = await regenerate_overview(LLMClient())
|
||
except Exception: # noqa: BLE001 — best-effort: never fail the 204
|
||
logger.exception("source removal: overview regeneration failed (best-effort)")
|
||
overview = False
|
||
bump_db = SessionLocal()
|
||
try:
|
||
bump_sources_version(bump_db)
|
||
bump_db.commit()
|
||
finally:
|
||
bump_db.close()
|
||
|
||
logger.info(
|
||
"source removed: kind=%s name=%s docs_pruned=%d files_removed=%s "
|
||
"overview=%s total_ms=%d",
|
||
row.kind,
|
||
name,
|
||
docs_pruned,
|
||
files_removed,
|
||
overview,
|
||
round((time.monotonic() - started) * 1000),
|
||
)
|
||
return Response(status_code=204)
|