34 Commits
Author SHA1 Message Date
arc53-machine 51f88144c5 Run webhooks with the owner's rules again
The owner sets a webhook up and its URL is a secret, and a webhook run
already refuses every action that needs approval, so holding it to the
API write allowlist only broke the owner's own automations. Webhook runs
no longer count as external callers; schedules keep their caller rules.
The allowlist copy no longer names webhooks.
2026-09-29 15:55:36 +01:00
arc53-machine daf6af57ae Keep outside callers' write limits in scheduled and webhook runs
A scheduled or webhook run acts as the agent's owner with no one to
approve, so a public-link user or API-key caller could have the agent
schedule a write and have it run on the owner's accounts. Runs now keep
the caller's rules: a schedule set by someone who reaches the agent only
by its public link runs as a public-link caller, one set through the API
(recorded as created_via 'api', migration 0042) and every webhook run as
an external caller, each with the agent's API write allowlist.
2026-09-29 15:44:24 +01:00
arc53-machine 121b6dd071 Search every agent source in scheduled and webhook runs
Headless runs searched only the agent's primary source, so an agent whose
knowledge sat in its extra sources answered a schedule or webhook without
it. They now take the primary and every extra source through the same
owner-or-sponsor check a chat uses, shared as one helper, and retrieve
through the per-source dispatcher so each source keeps its own settings.
2026-09-29 15:19:26 +01:00
arc53-machine 3511004f50 Connections own their credentials
Migration 0038 moves every stored secret (OAuth tokens, MCP OAuth
tokens and client registrations, API keys) into the connection's
encrypted envelope, links API-key tools to one connection per distinct
credential, allows several accounts per provider, and adds
credential_mode to sources and tools. OAuth MCP tools keep resolving
each member's own token, as they did before.

docsgpt.connectors.service is now the only reader of OAuth tokens:
get_valid_token_info refreshes under a row lock and persists rotated
refresh tokens, and a revoked grant flags the connection, pauses its
sources and notifies the owner. Loaders build from a connection
(BaseConnectorLoader.from_connection), so scheduled sync covers Drive,
SharePoint and Confluence sources with no browser. S3 and Reddit keys
stay on the connection instead of in remote_data.

New endpoints: POST /api/connections, /setup, /reconnect,
/picker-token, /claim, DELETE /api/connections/<id>, per-action
permissions and MCP refresh-tools. Upload, file listing, sync and
validate-session take a connection_id; session tokens keep working for
this release. The tool executor reads credentials from the resolved
connection (owner or member mode) and pauses on a Connect card when a
connection needs signing in. docsgpt connectors reencrypt rewrites
stored credentials after a key rotation.
2026-09-28 17:03:09 +01:00
arc53-machine 778830208b Write chat traces off the stream's thread
The OTel replay and the trace INSERT ran in the stream's finally, so a slow
database held the SSE connection open after the last event. The trace is
still frozen when the stream ends, but written on a small writer pool.
2026-09-23 22:42:05 +01:00
arc53-machine 5d0992eef8 Trace scheduled, webhook, search, MCP and graph-extraction runs
run_agent_headless records each unattended run under its endpoint; the
scheduler passes its run id and the webhook worker its task id so Logs rows
can find their trace, while the LLM's own request id stays untouched for
quota counts. /api/search and MCP search_docs record their retrieval, and a
graph build records every extraction call under one step.
2026-09-23 17:37:24 +01:00
arc53-machine d787f54c2d Keep the chunker's metadata when re-embedding a wiki page
`reembed_wiki_page_worker` built each chunk's metadata from scratch as
`{source, title, filename}`, discarding everything the chunker attached
to the document -- including the per-chunk `token_count` every other
ingest path records. Wiki chunks therefore reached the source viewer
with no count at all.

Carry `extra_info` over and let the page's own path and title win on top,
so a chunk that inherited a stale source from an earlier conversion is
still keyed and filtered correctly.
2026-09-22 13:44:24 +01:00
Alex 5406550bca feat(quotas): enforce user quotas on chat, agent, scheduled and webhook runs
check_usage now checks the billable user's quota on every request, before
the per-agent 24h limits, which keep applying to traffic through an agent.
Until now a request without an agent key skipped every limit. A refusal is
a 429 with Retry-After and a body naming the budget, usage, limit, the
layer the limit came from and when it resets.

Headless runs check the agent owner's quota before starting. A refused
scheduled run is recorded as budget_exceeded; a refused webhook run returns
a quota_exceeded result instead of raising, so Celery does not retry it.
2026-09-21 12:11:11 +01:00
Alex 3030aba39d fix: handle non text uploads more carefully 2026-09-14 17:27:13 +01:00
Alex 574f96341e refactor: rename the application package to docsgpt
The backend import package is now docsgpt, the name it will carry on PyPI;
application was far too generic to install into anyone's site-packages.
git mv plus a mechanical rewrite of every import, dotted string and path
reference: 734 Python files, the compose files, Dockerfile, workflows, docs,
setup scripts, devcontainer, k8s manifests, vscode config, pytest and coverage
config, .gitignore. Behaviour is unchanged.

Kept for one release:
- A top-level application package whose meta-path finder resolves
  application.x.y to the already-imported docsgpt.x.y object, so old imports
  and entry points (celery -A application.app.celery,
  uvicorn application.asgi:asgi_app) keep working with a FutureWarning.
- Celery registers every application.* task name as an alias of its
  docsgpt.* task on start-up, so messages queued by the previous release still
  run. The redbeat key prefix moves to redbeat:docsgpt:v2: so schedule entries
  the previous release wrote are left unread instead of firing twice.

The backend image builds from the repository root (docker build -f
docsgpt/Dockerfile .) so it can ship the alias package; a root .dockerignore
allow-lists docsgpt/ and application/ and keeps caches, local data, .env
files, the sample index files and the Dockerfile out. Compose and the image
workflows point at the new context.
2026-09-07 10:20:43 +01:00
Alex 9db11f18f5 fix(attachments,usage): validate content behind a BOM, judge by the live parser table
Review follow-ups.

A BOM told the sniff which encoding to read, but was also taken as the
verdict: three prepended bytes let any binary through, including as
notes.txt. A BOM now only selects the test — UTF-8 falls through to the
byte rules on the remainder, UTF-16/32 decode and judge the characters
(NUL, unprintable, or replacement chars from bytes the decoder could not
read). Real Notepad-Unicode text still passes, mp4-behind-a-BOM does not,
in either language.

The gate treated the full parser table as a given, but without docling the
fallback extractor has no .tif/.tiff/.bmp/.webp/.vtt/.xml handler, so those
suffixes skipped the content check and reached the plain-text fallthrough —
the original bug, one install away. The worker now passes the keys of the
extractor it actually built, making the second gate stricter than the
route's static one rather than a copy of it.

Cache bins: extraction coerced a missing value to 0 and only non-zero bins
were recorded, so a provider reporting cached_tokens=0 persisted as NULL —
indistinguishable from "not reported", and OpenAI reports exactly that on
every uncached request. Bins are now carried as Optional and recorded when
not None, which is what the nullable columns and the NULL-means-unknown
comment already assumed. Anthropic's cache_read/cache_creation bins had the
same shape and are fixed alongside; the int-or-None coercion is shared in
llm/base.py.
2026-09-03 09:46:55 +01:00
Alex 96cad1f8e2 fix(attachments): refuse unparseable chat attachments
A chat attachment with no parser fell through to SimpleDirectoryReader's
plain-text open(), so a phone-uploaded video was "extracted" into megabytes
of binary garbage, truncated, and stored with extraction.status == "ok".

Gate attachments in two tiers instead. A suffix with a dedicated parser is
admitted on its name — a PDF is binary and parses fine. Anything else has to
read as text: the first 8KB are sampled and refused on a NUL byte or too many
other control bytes. That keeps source, config and log files working through
the plain-text fallthrough, and keeps out videos, archives and renamed
binaries alike. The route checks the staged spool before anything is stored
or queued; the worker repeats the check where the local file exists, raising
the non-retryable AttachmentRejectedError.

SUPPORTED_ATTACHMENT_EXTENSIONS gains the parser-backed suffixes it was
missing (.tiff, .tif, .bmp, .webp, .vtt, .xml) and is now exactly the file
extractor's keys plus .txt, with a test asserting the two agree. The composer
applies the same rule client-side, so an unsupported file is named before it
costs an upload, and a test pins the frontend list to the backend one.

Attachment failures now show their reason inline under the chips rather than
only in a hover tooltip, which a touch user can never see, and only after a
send was attempted. Dropped `accept` from the dropzone: it discarded rejected
drops with no feedback and disagreed with the server about text files.
2026-09-03 08:52:12 +01:00
Alex 12dd7c7eab fix: stop the re-embed migration from destroying the index it rebuilds
Five defects from a review of the embeddings work, four of them silent.

- Write local files atomically. `LocalStorage.save_file` streamed straight onto
  the destination, so an interrupted write left a truncated file. `reembed`
  rewrites every index it touches, and a half-written `index.faiss` loads at
  neither the old width nor the new one -- the source was unrecoverable, with
  no backup and no temp file left behind. Bytes now land beside the destination
  and move into place with `os.replace`. S3 was already safe (single PUT).

- Read pgvector chunks a page at a time. `reembed_pgvector` materialised every
  `(id, text)` row for a source before embedding -- ~1.6 GB at 200k chunks and
  several times that for non-Latin scripts, with the `PGresult` held alongside
  until the cursor closed. Inside the shipped 4Gi limit, while also holding the
  model, that is an OOMKill -- which is exactly the SIGKILL the point above
  turned into a destroyed index. It now walks the source by keyset.

- Bound the first wave of delegated embeds. The failure cooldown is only latched
  once the first `get()` returns, so every request already in flight paid the
  full EMBEDDINGS_DELEGATE_TIMEOUT: measured 64 threads all timing out together,
  and at the shipped 60s across a 96-thread WSGI pool that is an API serving
  nothing at all, health checks included. One caller now probes while the rest
  fail fast; after a single success the gate leaves the path entirely.

- Ship EMBEDDINGS_NAME commented in .env-template. The comment directly above it
  says to leave it commented when upgrading, and the line shipped set. Any value
  reaching `.env` lands in `model_fields_set`, which makes `resolve_embeddings_pin`
  bail -- so a template-derived `.env` disabled the legacy pin outright and
  repointed a populated index at a different 768-dim model, where no width check
  fires. The pin already picks granite for a fresh install and mpnet for an
  existing one, so nothing needs to be set by hand.

- Stamp `sources.model` on wiki sources. They were created with the column NULL
  and then embedded like any other source, and the boot check reads NULL as
  "pre-dates the column, therefore the legacy model" -- reporting a correctly
  embedded source as stale on every startup of every process. Stamped at
  creation, and again on each page re-embed so existing rows heal.

The two docs that promised the FAISS index survives a failed run said so of the
embed only; both now describe the write, and upgrading.mdx says to stop ingest
for the duration.
2026-08-28 16:08:15 +01:00
Alex e0aff39a1b feat: ingestion optimisations 2026-08-25 12:27:09 +01:00
Alex 7ad2f38512 feat: faster attachment processing 2026-08-13 16:04:45 +01:00
Alex 0a15ce8fbb feat: attachment provenance 2026-08-13 14:30:28 +01:00
Alex 0ecb421954 fix: error type fixes and docling parsing improvements 2026-08-06 12:29:51 +01:00
Alex 15b6b03b47 fix: minor size cleanup 2026-08-04 15:07:54 +01:00
Alex a1ad802e22 feat: limit docling attachment file sizes explicitly 2026-08-04 14:04:44 +01:00
Alex 94a845aa82 fix: more artefact hardening 2026-07-04 11:42:27 +02:00
Alex 37d93cbd86 Parse documents on a Celery parsing worker via a read_document tool
Replace the sandbox Docling extractor with read_document, backed by the in-process
backend parser (the same one ingestion uses) and offloaded to a dedicated
'parsing' Celery queue so it can run on GPU-capable workers with predictable RAM.
The tool resolves the input ref under the run-scoped gate, enqueues the parse,
and awaits it with a timeout (degrading to an error rather than hanging); the
worker independently re-resolves the artifact through the same gate and never
trusts a raw path. Untrusted files get the upload path's safeguards (extension
whitelist, size cap, sanitized temp file, cleanup). Options: output
(markdown/text/structured/chunks), ocr, pages, engine, max_chars, include_tables,
persist, json_schema. The workflow native-file 'extract' fallback now uses the
same worker path, so document parsing no longer needs the sandbox and works on
every backend.

Also fixes the branch's periodic-task test (the sandbox reaper made it 12) and
points the dev and e2e Celery workers at the parsing queue.
2026-06-25 13:24:12 +01:00
Alex d7bbfcfe17 fix: minor graph rag improvements 2026-06-23 20:10:41 +01:00
Alex 4742aec4c6 feat(graphrag): durable extract_graph task + ingest-path enqueue + enable route
graph_enabled() sets kind=graphrag + retriever=graphrag. extract_graph_worker fetches the source's pgvector chunks and runs G3 extraction (graphrag_available guard, empty no-op). extract_graph durable+idempotent task; key varies with source updated_at so re-ingest/re-enable re-run incrementally (G3 checkpoint skips done chunks) while concurrent same-state enqueues dedup. The 4 ingest paths enqueue after embed when kind=graphrag (isolated in try/except so a broker hiccup can't fail the ingest). POST /api/sources/<id>/graphrag/enable: pgvector+GRAPHRAG_ENABLED gated, owner/editor write-authz; PATCH config still can't flip kind->graphrag. Unit G4.
2026-06-23 01:29:51 +01:00
Alex ef44459984 fix: small wiki fixes 2026-06-23 00:19:56 +01:00
Alex 8a0f7024e3 fix(wiki): convert reassembles pages from existing chunks (crawler/remote)
convert_source_to_wiki now builds pages from the source's already-ingested vector-store chunks instead of re-parsing storage files, so it works for crawler/remote/connector sources (which have no stored files; file_path empty). Groups by metadata file_path/file_name (URL- and connector-aware; URLs normalized to virtual paths) rather than the raw source. Passes embeddings_key to create_vectorstore in BOTH convert and reembed_wiki_page (fixes a TypeError crash on faiss/elasticsearch). Conservative overlap-trim (min length). Deletes the original chunks after reassembling so retrieval has no duplicate/stale content.
2026-06-22 23:03:52 +01:00
Alex 0c9e0313bd feat(wiki): convert existing source to wiki + human page edit endpoint
Unit 5 of F-Wiki (D20-D23). convert_source_to_wiki task reuses reingest's storage file-load + parser to materialize files->wiki_pages (one page/file), skips/reports non-text, re-embeds per page, and flips kind=wiki + exposure=agentic_tool only when pages were created. POST /wiki/convert (explicit, write-authz; blank source enables inline, fileful enqueues the task; rejects mid-ingest). PUT /wiki/page for human edits (write-authz, optimistic version -> 409, re-embed). PATCH /config preserves kind (kind changes only via convert).
2026-06-22 20:16:18 +01:00
Alex d68be86244 feat(wiki): reembed_wiki_page durable task
Per-page re-embed (Unit 2 of F-Wiki): targeted delete of the page's old chunks, re-chunk via the source's chunking config, add_chunk with reingest-matching metadata (source=path), set embed_status embedded/failed. Durable + idempotent (key=content_hash), mirroring reingest_source_task. Missing page => purge only.
2026-06-22 17:49:04 +01:00
Alex f6400cd736 feat: per-source RAG configuration (retrieval strategies, chunking, exposure, prescreen)
Introduces a per-source config contract that makes RAG behavior strategy-dispatched instead of a single hardcoded path. Every source gains a validated JSONB config; an empty/absent config reproduces current behavior byte-for-byte, and the whole path is gated by PER_SOURCE_RETRIEVAL_ENABLED.

Foundation: sources.config JSONB column + migration 0022_source_config; SourceConfig/ChunkingConfig/RetrievalConfig pydantic models (strict on write, lenient on read); ChunkerCreator and RetrieverCreator.register registries; config threaded through the upload routes, ingest/remote/connector workers, and reingest.

Retrieval: a Dispatcher groups sources by retriever key (all-classic collapses to today's single ClassicRAG under one shared token budget; non-classic retrievers get their own instance), removing the previous single-global-retriever collapse in stream_processor. Per-source chunks, score_threshold (honored for pgvector/mongodb, safely ignored elsewhere), and rephrase_query toggle. New PATCH /api/sources/<id>/config with team-aware (effective_write_owner) authz and a requires_reingest signal.

Chunking strategies: recursive, markdown, parent_child (selectable per source; re-ingest to apply). Search exposure: per-source prefetch vs agentic_tool for agentic/research agents. Map-reduce prescreen: optional LLM relevance pre-filter implemented as a composable post-retrieval stage that wraps any retriever.

Backend and frontend (shared Retrieval options panel + edit modal) with tests; backend suite and frontend vitest green. Excludes the wiki and GraphRAG flagships.
2026-06-20 21:54:23 +01:00
Alex 0507b32221 feat: default tools (#2485)
* feat: default tools

* fix: tests

* feat scheduler

* fix: scheduler UI

* Agent switch breadcrumbs

* fix: minor fixes to scheduler

* fix: tests
2026-05-22 16:05:03 +01:00
Alex e167cf8247 fix: broken syncs (#2480)
* fix: broken syncs

* fix: mini fixes
2026-05-17 23:58:28 +01:00
Alex e351f45d88 Feat notification system (#2472)
* feat: SSE notification system

Adds a per-user SSE pipe (GET /api/events) plus a per-message
chat-stream reconnect endpoint (GET /api/messages/<id>/events).

Backend substrate:
- application/events/ — durable journal (Redis Streams) + live
  pub/sub for user-scoped events, with publish_user_event() as
  the worker-side entrypoint.
- application/streaming/ — broadcast_channel for pub/sub fanout
  and event_replay for the per-message snapshot+tail path.
- application/storage/db/repositories/message_events.py +
  alembic 0007 — Postgres journal for chat-stream events.
- application/worker.py — ingest/reingest/remote/connector/
  attachment/mcp_oauth tasks publish queued/progress/completed/
  failed envelopes alongside their existing status updates.

Frontend client:
- frontend/src/events/ — connect/reconnect, Last-Event-ID cursor,
  backoff with jitter. Each tab runs its own connection; no
  cross-tab dedup (future work).
- frontend/src/notifications/ — recentEvents ring, cursor
  tracking, tool-approval toast.
- frontend/src/upload/uploadSlice.ts — extraReducers for
  source.ingest.* and attachment.* events.

Coverage: 132 SSE tests across events substrate, replay, journal,
routes, and worker publishes.

* refactor(attachments): remove polling, SSE-only

frontend/src/components/MessageInput.tsx no longer runs a 2s
setInterval against getTaskStatus for every processing
attachment. The attachment.* SSE reducers in uploadSlice.ts are
now the sole driver of attachment state transitions.

* feat(connector): consume source.ingest.* SSE, remove polling

frontend/src/components/ConnectorTree.tsx now mirrors FileTree's
slice-walking pattern: it watches notifications.recentEvents
for source.ingest.{completed,failed} envelopes matching the
sync's source id, and no longer polls /task_status every 2s.

* refactor(source-ingest): remove polling, SSE-only

frontend/src/upload/Upload.tsx and
frontend/src/components/FileTree.tsx no longer run getTaskStatus
polling fallbacks. The source.ingest.* SSE reducers in
uploadSlice.ts and FileTree's slice walk are now the sole
drivers of upload/reingest state transitions.

* refactor(mcp-oauth): carry authorization_url in SSE, remove polling

application/worker.py::mcp_oauth now publishes
authorization_url on the mcp.oauth.awaiting_redirect envelope.
frontend/src/modals/MCPServerModal.tsx consumes it from SSE
instead of polling /oauth_status/<task_id> every 1s.

The URL is generated inside DocsGPTOAuth.redirect_handler when
the FastMCP client triggers OAuth. The worker now plumbs a
publish callback through tool_config -> MCPTool -> DocsGPTOAuth
so the awaiting_redirect publish fires from inside the handler
at the exact point the URL becomes known. The legacy Redis
mcp_oauth_status setex writes and the GET
/api/mcp_server/oauth_status/<task_id> endpoint are kept as
belt-and-suspenders; nothing in the frontend reads them now.

* feat(source-ingest): plumb limited flag through SSE for token-cap UX

application/worker.py::ingest_worker and remote_worker now publish
``limited: bool`` on the source.ingest.completed envelope.
uploadSlice routes ``payload.limited === true`` to a failed status
with a ``tokenLimitReached`` flag, and UploadToast surfaces the
translated tokenLimit i18n string. No worker code path sets
limited=true today; this is a forward-looking contract so when
token-cap detection lands, the UX is already wired.

* refactor(mcp-oauth): read status from SSE journal, drop polling endpoint

MCPOAuthManager.get_oauth_status now walks the per-user SSE Streams
journal (user:{user_id}:stream) for the latest mcp.oauth.* envelope
matching the task id, returning the status string derived from the
event type suffix and the payload fields. The worker is the single
source of truth — its publish_user_event calls write the same
record the SSE client receives live.

Removed:
- /api/mcp_server/oauth_status/<task_id> route in
  application/api/user/tools/mcp.py
- mcp_oauth_status worker function and mcp_oauth_status_task Celery
  wrapper
- All mcp_oauth_status:{task_id} Redis setex writes (4 in mcp_oauth,
  2 in DocsGPTOAuth.redirect_handler / callback_handler)
- The update_status closure in mcp_oauth that wrote the polling
  payload

Tests updated:
- get_oauth_status now takes (task_id, user_id); new coverage walks
  a fake xrevrange response for the completed envelope, the no-match
  case, and a Redis-down case
- Removed TestMCPOAuthStatus route tests and TestMcpOauthStatusTask
  celery-wrapper test
- Removed the two oauth_status methods from the integration runner

mcp_oauth:auth_url/state/code/error Redis keys remain — they are
the OAuth flow's own state (not the dropped polling payload).

* chore(mcp-oauth): delete orphaned getMCPOAuthStatus client

The /api/mcp_server/oauth_status/<task_id> endpoint was removed in
the prior commit; the corresponding userService method and the
MCP_OAUTH_STATUS endpoint constant had no remaining callers in the
frontend, so they're deleted along with it.

* fix(events): drop live publish when journal write fails

application/events/publisher.py returned an envelope to live
pubsub subscribers even when the XADD to the durable journal
failed. The envelope had no ``id`` field, which bypassed the SSE
route's dedup floor and broke ``Last-Event-ID`` semantics for any
reconnecting client.

Best-effort delivery means dropping consistently, not delivering
inconsistent state. Now: if the journal write fails the publisher
returns None and skips the live publish entirely.

* fix(notifications): dedupe sseEventReceived against immediate dupes

Snapshot replay + live tail can both deliver the same id when the
live pubsub frame and the replay XRANGE overlap. The route's own
dedup floor catches the common case, but consumers walking
``recentEvents`` (FileTree, ConnectorTree, MCPServerModal,
ToolApprovalToast) would otherwise act on the same envelope
twice when a duplicate slipped through.

Belt-and-suspenders: short-circuit when the most recent id in
the ring matches the incoming one.

* fix(events): skip replay budget INCR when no snapshot work possible

_allow_replay incremented the per-user counter on every
/api/events GET, including no-op connects from a fresh client
with no cursor against an empty backlog. React StrictMode dev
double-mounts plus a few tabs trivially tripped the default
30-per-60s budget on idle reconnects.

XLEN pre-check: when last_event_id is None and the user stream
is empty, the connect can't do snapshot work — return True
without INCR. Cursor-bearing connects still INCR unconditionally
(probing the cursor's relationship to stream contents would
require a redundant XRANGE).

* fix(streaming): tighten journal contract + recover from seq collisions

Two related fixes to application/streaming/message_journal.py.

1. record_event now rejects non-dict payloads at the gate. The
   live path (base.py::_emit) wrapped non-dicts as
   {"value": payload}; the replay path in event_replay synthesized
   {"type": event_type}. A reconnecting client would receive a
   different envelope than the one originally streamed. Now both
   paths see byte-identical envelopes because non-dicts can't be
   journaled at all. The corresponding event_replay fallback is
   replaced with a warn-and-skip for any legacy rows.

2. record_event handles IntegrityError on (message_id, sequence_no)
   collisions by reading latest_sequence_no and retrying once with
   latest+1. The most likely cause is a stale seq seed on a
   continuation retry where the route read MAX(seq) from a
   separate connection before another writer committed past it.
   Previously the error was swallowed and the event silently
   dropped from the journal; now it lands at the next available
   seq. The live pubsub publish uses the materialised seq so the
   journal row and the live frame agree.

* perf(streaming): batch message_events INSERTs per stream

complete_stream previously opened a fresh db_session() per yielded
event, doing one Postgres INSERT + commit per chunk on the WSGI
thread. Streaming answers emit ~100s of answer chunks per response,
so the route was paying ~100 PG roundtrips per stream serialized on
commit latency.

New BatchedJournalWriter in application/streaming/message_journal.py
accumulates rows per stream and flushes on three triggers:
- size: buffer reaches 16 entries
- time: 100ms elapsed since the last flush
- lifecycle: close() at end-of-stream

Live pubsub publishes still fire synchronously per record(), so
subscribers see events in real time — only the durable journal write
is amortized. On bulk INSERT IntegrityError the writer falls back to
per-row record() with the existing seq+1 retry so a single colliding
seq doesn't drop the rest of the batch.

complete_stream wires journal_writer.close() into every exit path
(happy end, tool-approval-paused end, GeneratorExit, error handler)
so the terminal event is committed before the generator returns —
otherwise a reconnecting client could snapshot up to the last flush
boundary and live-tail waiting for an end that's still in memory.

Repository gets bulk_record() — one SQLAlchemy executemany INSERT
for the bulk path. All-or-nothing on collision (Postgres aborts the
whole batch); the writer's per-row fallback handles recovery.

* chore(upload): drop dead UploadTask.lastEventAt field

The lastEventAt field on UploadTask had no remaining consumers — the
matching Attachment.lastEventAt was cleaned up earlier. Remove the
field declaration and the slice write site.

* chore(frontend): drop orphaned getTaskStatus client

After the polling-removal sweep no caller in frontend/src/ references
userService.getTaskStatus or endpoints.USER.TASK_STATUS. The backend
route /api/task_status itself stays — agents, webhooks, e2e specs,
and the public docs still depend on it.

* docs(repo): remove stale planning docs from repo root

notification-channel-design.md, plan.md, and reminder-tool-design.md
were leftover Claude planning artifacts from the SSE substrate work
that landed accidentally. CLAUDE.md prohibits creating planning docs
unless asked — delete them.

* docs(message-events): clarify repo vs wrapper payload contract

MessageEventsRepository.record accepts any JSONB-compatible value; the
streaming wrapper record_event tightens this to dicts only because the
live and replay paths reconstruct non-dict payloads differently. Spell
the split out so the next reader of the repo method doesn't assume the
wrapper's contract applies here.

* refactor(events): raise on malformed stream id instead of lex fallback

stream_id_compare's lex-fallback branch was a footgun: a malformed id
that sorts lex-greater than a real one would pin live-tail dedup
forever, dropping every subsequent legitimate event silently. Both
current callers in application/api/events/routes.py pre-validate
inputs against _STREAM_ID_RE before calling, so changing the function
to raise ValueError is a no-op on the happy path and turns the future-
caller footgun into a loud failure.

* test(tasks): cover cleanup_message_events task body

Adds skipped-when-no-POSTGRES_URI and happy-path coverage for the
Celery janitor. The skipped path returns the documented short-circuit
shape without touching the repo. The happy path seeds a backdated
row, runs the task against the pg_conn fixture, and asserts the
retention window's row is deleted while in-window rows survive.
Mirrors the TestCleanupPendingToolState pattern.

* fix(notifications): treat /c/new as no current conversation

useMatch('/c/:conversationId') treats the literal URL /c/new as a
real conversation id, so the toast suppression check confused
'user is on /c/new' with 'user is on the conversation needing
approval'. Explicit guard: when the matched id is 'new', fall
through to the no-match case so approval toasts still surface.

* docs(events): enumerate publish_user_event None-return paths

The function returns Optional[str] today, with None conflating five
distinct outcomes (missing args / push disabled / unserialisable /
Redis down / XADD failed). Every current call site is fire-and-
forget and ignores the return, so the right move is to document the
five cases rather than promote to an enum return — keeps the API
small while making the diagnostic surface (logs) obvious. If a
future caller needs to react differently per reason, promote then.

* refactor(sources): move source-id derivation out of worker module

application/api/user/sources/upload.py imported _derive_source_id
from application.worker — pulling the entire Celery worker module
into the API process at import time just for a two-line helper.

Move DOCSGPT_INGEST_NAMESPACE and the derivation function to a
new application/storage/db/source_ids.py module that both layers
can import without that dependency edge. worker.py re-exports the
old names (_derive_source_id, DOCSGPT_INGEST_NAMESPACE) for
backward-compatible imports from tests and any other in-tree
callers; new code should import from the new module directly.

* fix(cache): enable Redis health_check_interval to surface half-open TCP

Without health_check_interval, a half-open TCP socket (NAT silently
dropped state, ELB idle-close) can leave pubsub.get_message hanging
past the SSE generator's keepalive cadence — the kernel never
surfaces the dead socket because no payload is in flight. Setting
health_check_interval=10 makes redis-py ping every 10s when
otherwise idle, so the next get_message after the dead window
raises and the SSE loop falls into its reconnect path instead of
silently freezing on the user.

* chore(events): rename attachment.processing.progress to attachment.progress

The event-type taxonomy was inconsistent: source ingest emits
source.ingest.progress (three segments) while attachments emitted
attachment.processing.progress (four segments). Drops the
.processing. infix for parity. Worker publish sites, the slice
reducer's match, and the worker tests all flip together.

No external consumers — the event type is purely internal between
the publisher and the in-tab slice; safe to rename in one commit.

* feat: events cleanup

* fix: better docs

* fix: e2e tests
2026-05-15 12:23:31 +01:00
Alex b4c4ab68f0 feat: durability and idempotency keys (#2450)
* feat: durability and idempotency keys

* feat: more durable frontend

* fix: tests

* fix: mini issues

* fix: better json validation

* fix: tests
2026-05-04 23:25:41 +01:00
Alex 318de18d43 feat: BYOM (#2433) 2026-04-27 22:09:33 +01:00
81b6ee5daa Pg 4 (#2390)
* feat: postgres tests

* feat: mongo cutoff

* feat: mongo cutoff

* feat: adjust docs and compose files

* fix: mini code mongo removals

* fix: tests and k8s mongo stuff

* feat: test fixes

* fix: ruff

* fix: vale

* Potential fix for pull request finding 'CodeQL / Clear-text logging of sensitive information'

Co-authored-by: Copilot Autofix powered by AI <62310815+github-advanced-security[bot]@users.noreply.github.com>

* fix: mini suggestions

* vale lint fix 2

* fix: codeql columns thing

* fix: test mongo

* fix: tests coverage

* feat: better tests 4

* feat: more tests

* feat: decent coverage

* fix: ruff fixes

* fix: remove mongo mock

* feat: enhance workflow engine and API routes; add document retrieval and source handling

* feat: e2e tests

* fix: mcp, mongo and more

* fix: mini codeql warning

* fix: agent chunk view

* fix: mini issues

* fix: more pg fixes

* feat: postgres prep on start

* feat: qa tests

* fix: mini improvements

* fix: tests

---------

Co-authored-by: Copilot Autofix powered by AI <62310815+github-advanced-security[bot]@users.noreply.github.com>
Co-authored-by: Siddhant Rai <siddhant.rai.5686@gmail.com>
2026-04-18 13:13:57 +01:00