582 lines
25 KiB
Python
582 lines
25 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``), ``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), ``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). 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, import, or prune anything — the existing Sync button performs
|
||
that (a removal prunes on the next sync, ``prune=True``). The upload
|
||
route is the 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 GitSource
|
||
from app.rag.archive_upload import (
|
||
ARCHIVE_SUFFIXES,
|
||
ArchiveUploadError,
|
||
archive_source_name,
|
||
swap_in,
|
||
unpack_archive,
|
||
)
|
||
from app.rag.importer import import_sources
|
||
from app.rag.llm import LLMClient, check_models
|
||
from app.rag.overview import regenerate_overview
|
||
from app.schemas import (
|
||
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@)")
|
||
|
||
|
||
@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`` 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`` 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,
|
||
)
|
||
for row in rows
|
||
],
|
||
from_env=False,
|
||
)
|
||
return GitSourceList(
|
||
sources=[
|
||
GitSourceRow(id=None, kind="git", url=url, path=None, added_at=None)
|
||
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.
|
||
"""
|
||
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)
|
||
|
||
|
||
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"), "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),
|
||
f"a local source with this path already exists: {path}",
|
||
db,
|
||
)
|
||
|
||
|
||
@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>``);
|
||
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>)``
|
||
+ 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:
|
||
if db.scalar(select(GitSource).where(GitSource.path == path)) is None:
|
||
db.add(GitSource(url=path, kind="local", path=path))
|
||
try:
|
||
db.commit()
|
||
except IntegrityError:
|
||
db.rollback()
|
||
raise ValueError(
|
||
f"a local source with this path already exists: {path}"
|
||
) from None
|
||
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)
|
||
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)
|
||
def delete_git_source(
|
||
source_id: uuid.UUID,
|
||
db: Session = Depends(get_db), # noqa: B008
|
||
) -> Response:
|
||
"""Remove a stored row; 404 when the id is unknown.
|
||
|
||
Removing does not touch the clones or the index — the next Sync
|
||
(``prune=True``) prunes the dropped repo (phase scope boundary).
|
||
"""
|
||
row = db.get(GitSource, source_id)
|
||
if row is None:
|
||
raise HTTPException(status_code=404, detail="git source not found")
|
||
db.delete(row)
|
||
db.commit()
|
||
return Response(status_code=204)
|