A paused turn is found by the agent owner's id, so anyone holding one of
the owner's agent keys could resume the owner's own chat with its saved
wiki edit rights. A resume now counts as an API or widget caller when
either the saved state or the resuming request is one, cuts the wiki tool
to wiki_view unless the wiki allows outside edits, and gives the tool
executor the same flags. A request that names an agent, by key or id, may
only resume that agent's turn; otherwise the claim is released and the
request refused.
Public-link visitors run as themselves and reach only wikis they may edit,
so the wiki switch no longer applies to them. Instead every wiki write in
a public-link run waits for the visitor's approval, so the agent owner's
prompt or sources can't steer an edit to the visitor's wiki unasked.
A run from an agent's API key or widget acts as the agent's owner, so it
could rewrite any wiki the owner can edit. A new per-wiki setting,
wiki_outside_edits (off by default), decides whether such runs, and runs
from the agent's public link, get the wiki's edit actions. While it is off
they are offered only wiki_view, and the tool refuses writes itself after
reading the live setting. The owner changes it through the owner-only
/api/sources/<id>/wiki/settings route; tokens can read it but not change it.
The answer routes refuse such a request with 401, but only after building
the agent, so a public agent's prompt tools were pre-fetched, and their
actions run, for nobody. Anonymous chat without an agent key is not
supported, so the processor now stops before any setup and the route
answers 401 as before.
Tool pre-fetch ran the caller's own active tools, so a teammate chatting
with a shared agent had the owner's prompt fill in from their tools, and
even the owner got every active tool rather than the agent's. It now uses
the toolset the run gets: the agent's tools as its owner or sponsor, or
the caller's tools and defaults outside an agent. Pre-fetch asks nobody,
so on someone else's tool it skips approval-gated actions and anything on
a connected account.
chunks is a total per request, split across the attached sources, so an
agent with two sources and the default of 2 got a single chunk from each,
and its answers changed with whichever chunk won. 6 gives three per
source for about 4-5k more input tokens per retrieval turn.
Every literal default moves from 2 to 6: the request default, the
retrievers, the internal search tool, workflow agent nodes, scheduled and
headless runs, agent create/update/import, the source retrieval config and
the frontend forms. Existing agents and sources keep what they store; a
source saved with chunks=2 now counts as configured at 2, which is pinned
by a test.
Headless runs also read chunks=0 as unset (`or 2`), so an agent with
retrieval switched off retrieved anyway on scheduled runs; 0 now stays 0.
A turn whose agent raised wrote no user_logs row, so it only surfaced as
the agent's system error row. Every finished turn now writes its chat row,
at level error with the error when it failed, and linked to its trace; the
system row for the same traced activity is no longer listed twice.
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.
build_agent no longer takes request_id from the request body: it becomes
the primary LLM's usage request id, and quotas count distinct request ids,
so a client could make every call count as one. Requests refused after
setup started (unauthorized, over quota, resume conflict, setup error) now
write their trace, marked error, through an after-request hook; streaming
routes hand the trace to complete_stream instead.
Tool results, denial comments and tool exception text now reach a span only
as a capture-gated preview; span.error is a fixed message, since it is
stored and exported regardless of content settings. A turn whose stream
yields an error event is recorded as failed. Search traces are listed by a
query copied into the small summary column, so the Logs timeline never
reads the spans JSONB. Adds GraphRAG span tests.
StreamProcessor starts the trace and mints the request id before the
agent is built, so pre-fetch retrieval and compression are inside it and
side-channel LLM calls share the id. complete_stream activates it in the
SSE pump thread, binds the message and conversation, and writes it once
however the stream ends: paused, failed, abandoned or superseded (dropped).
user_logs rows now carry request_id and message_id.
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.
An unchained request with no system message staged None, and the record
step skipped the commit, so the previous head hash survived a transcript
that never received a head; a later chained request restoring that head
would have omitted it. The staged value is now committed as-is, None
included. Also pins each rejection predicate of is_usable_compression_point
with its own test.
Three review findings on the bounded-chain change.
Mid-execution compression rebuilt the conversation from the in-flight
messages, which after a turn-start reuse hold only the recent turns: the
summary living in the system prompt never reached the compressor, so the
new summary replaced the old one, and the persisted point's query_index
was relative to that shortened list. The summary the agent is running
under now rides into the synthetic conversation as its latest point
(query_index -1, so every in-flight query is new), for both the database
and the in-memory path, and the database path persists the index of the
saved conversation's last row.
Saved points with an empty summary, which earlier versions wrote, were
treated as reusable: get_compressed_context sliced the raw history away
and the effective token count made the conversation look small. Point
selection everywhere now takes the latest usable point (non-blank
summary, positive token count) and falls back to the raw history when
there is none.
The chained system-head hash was committed while building the request,
so a transport failure followed by the same-primary retry omitted a
changed system message. The hash is now staged per request and committed
only when the provider records the response.
- append_compression_point only skips a point when both query_index and
compressed_summary are present and match the last one; points without
those fields (as in the repository tests) were all being treated as
duplicates.
- The incremental compression tail and the orchestrator's "anything new
since the last point" check exclude the visible summary row, which the
prompt already receives through existing_compressions.
- Summary rows carry a persisted metadata marker; replay filters on the
marker, and falls back to the label only for rows written before it that
have no tool calls and no per-turn metadata, so a user who types the
label text keeps their turn.
- The prompt_cache_key is a hash of the user id, never the id itself.
- Describe truncation="auto" as dropping the oldest items.
In store mode every user turn chained onto the previous response, so the
provider's stored transcript grew without bound (measured: 889k prompt
tokens for a 37k-token saved history) while every local guard, the
compression pipeline included, measured the saved history. Each chained
tool round also re-sent the system message, which the server appends rather
than dedupes, and a saved compression point was applied exactly once, in the
turn that made it.
Chaining is now bounded. A turn starts from the saved history when the
previous turn's reported prompt reached the chain budget (default: the
model's context window), when the conversation was compressed after that
turn was produced, or when OPENAI_RESPONSES_CHAIN_ACROSS_TURNS is off.
Chained rounds omit an unchanged system head (hash carried in the persisted
Responses state). truncation="auto" is available behind a setting as a
backstop against a chain that outgrows the model's window.
Compression: a saved point is applied at every turn start; the threshold
counts the summary plus the queries after the point instead of the raw
history; re-compression summarises only the tail on top of the last point;
the mid-execution path marks itself persisted and resets the provider chain
so the rebuilt messages are the context; an empty summary is rejected; the
visible "[Context Compression Summary]" rows are no longer replayed as
history; appending the same point twice is a no-op.
Cache hints: a per-user prompt_cache_key and an optional
prompt_cache_retention on Responses API calls.
Measured on Azure with the same client shape as production (stateless
OpenAI client, server-side tools, PDF part): tokens billed on the sixth turn
fell from 58k to 35k, tool rounds add tens of tokens instead of ~2.8k, the
turn after a compression reused the saved summary in under two seconds
instead of re-summarising, and the round after a mid-execution compression
started from the compressed context instead of the full stored transcript.
Source access control
---------------------
`active_docs` is client-supplied and reached the retriever unchecked, and the
retriever queries `WHERE source_id = <id>` with no owner predicate — so any
caller could pass any source id to /stream or /api/answer and have another
tenant's documents quoted back, while /api/sources/<id>/search correctly
refused the same id. Gate it through `can_access`, the helper the guarded
endpoints already use, and filter `self.source` down to the authorized set.
Fails closed: no principal, or a check that errors, drops the source.
Three sibling paths had the same gap:
- workflow agent nodes: `AgentNodeConfig.sources` is written verbatim from
client JSON at save time and nothing validated it, so a node could name any
tenant's source. Gate against the workflow owner, so shared workflows keep
reading their owner's sources like shared agents do.
- /api/share: `_resolve_source_pg_id` resolved any id with no ownership
predicate and baked it into the agent the share creates; /api/search then
searched it. Authorize before attaching.
- search_service: re-resolve the ids stored on an agent row instead of
trusting them, so a row written by any future path with the same gap cannot
be read back.
Team grantees previously lost their source's retrieval config: the post-check
read was still owner-scoped, so it missed and fell back to defaults (an
`agentic_tool` source was bulk-prefetched for every grantee). Read unscoped
after `can_access` passes.
Retrieval
---------
`PGVectorStore._ensure_table_exists` created an IVFFlat index on the empty
table it had just created. IVFFlat computes centroids at build time, so those
centroids were random, and combined with the `source_id` post-filter a source
with hundreds of embedded chunks returned zero rows — retrieval reported no
documents, the model answered from memory, and nothing was logged. Stop
creating the index (exact search is correct and fast well past the sizes most
deployments reach); raise `ivfflat.probes` to sqrt(lists) where an index still
exists; and re-run a short indexed search exactly, since post-filtering means
no index setting can guarantee a full result. `graphrag` had the same
empty-table index with no fallback at all.
Also: bound `chunks` to 0-500 on both the request and agent paths (0 still
means "skip retrieval"), let a source's configured `retrieval.chunks` outrank
the request body, and cap ClassicRAG's per-source floor at
max(top_k, n_sources) so attaching sources cannot inflate the result set.
Silent failures
---------------
An empty retrieval was invisible to both the model and the client: the `source`
event was suppressed when the list was empty, so "searched and found nothing"
looked identical to "no source attached", and the prompt said nothing at all.
Emit the event always, and tell the model when a search ran and returned
nothing. A file that parses to nothing now fails ingest with a message naming
the cause instead of storing an embedding of the empty string. `score_threshold`
returns warnings when the active store or retriever cannot honour it.
Prompt structure
----------------
Retrieved documents move from the system prompt into the user turn, with the
injection guard restated next to them: they change every turn (defeating prefix
caching), they are third-party text that should not carry system authority, and
routing them through the query budget makes them truncatable rather than
silently crowding it out. Documents are shed lowest-ranked-first before the
question is touched.
The six chat presets (3 tones x 2 retrieval modes) differed only in their
Answering section; they are now composed from single-source fragments at load
time, not through Jinja inheritance, which would have opened a file-read
surface in the template sandbox and broken the tool-prefetch parser. Per-tool
guidance moves out of the prompt into tool schemas, so it travels with the tool
and cannot render when the tool is absent. A plain-text custom prompt is staged
as a persona value inside the skeleton instead of replacing it wholesale — it
used to silently lose the injection guard, platform block, memory and
attachments, and its braces are now inert.
Other fixes
-----------
- agents/base: an oversized system prompt drove the query budget negative and
dispatched a full-price request with an empty question; raise instead.
- llm/anthropic: migrate off the retired Text Completions API. It flattened
history to first+last message and ignored tools entirely. Adds the missing
Anthropic handler, without which every tool call was silently dropped.
- sources/upload: `sitemap` had no branch, so every sitemap ingest died on a
TypeError; `validate_url` now rejects a falsy URL cleanly.
- workflow nodes: retrieved documents never reached the node agent, so a
classic node with a source and an ordinary prompt answered "I have no
documents" while the run reported completed.
- parser/bulk: copy the metadata dict, or every chunk reports the last chunk's
token_count.
- crawler_loader: carry the page title, or citations render the whole chunk
body as the label.
WorkflowEngine reports node failures by *yielding* `{"type": "error"}`
rather than raising, so complete_stream's generator returns normally and
the except handler never runs. The turn was finalized `status="complete"`
with an empty response.
Live, the client renders an error bubble with a Retry button. On reload
it does not: mapServerQueryToClient only surfaces `metadata.error` for
`failed` rows, so history showed a blank message with no error and no way
to retry. A user hitting this re-sends the same prompt into new
conversations, which is exactly what the 2026-08-01 report shows — nine
blank first messages in seven hours.
Tracks a `stream_error` flag alongside the existing `paused` machinery
and finalizes `failed` when the turn produced no answer, recording the
user-facing message in `metadata.error`. An error arriving *after* output
keeps `complete` so partial text is not discarded; structured answers
count as output too, since they live in `structured_chunks` rather than
`response_full`. The flag is recorded before the pause branches so those
paths cannot lose it.
save_conversation grew a `status` parameter (default `complete`) for the
non-WAL branch, which took no status and so landed on the column default
— the same blank-complete row on a path the WAL fix did not cover.
Title generation now also runs for failed turns: _maybe_generate_title
only regenerates while the name is still the question-prefix fallback, so
skipping it would strand a conversation whose first turn failed with the
raw prompt as its name forever.
logging.py counts a yielded error toward `activity_finished.status`.
These failures previously logged `status=ok` with `answer_length=0`,
which is why user-visible blank answers never appeared in error metrics.
Token accounting:
- Drain each tool round's provider stream to exhaustion before running
tools and recursing, so the usage decorator persists exactly one
token_usage row per LLM call, at call end. Previously every round's
generator was abandoned mid-iteration and flushed together at request
teardown, writing N near-identical rows stamped with the final
round's provider counts (duplicate billing).
- Consume the Chat Completions include_usage terminal chunk (it arrives
after finish_reason and was never read) so streamed calls record
provider-exact token counts instead of tiktoken estimates.
- Claim provider-reported usage once per call (_last_usage_claimed) so
a late-finalized generator can never adopt another call's counts.
Oversized-context guards:
- Enforce Responses API function_call/function_call_output pairing in
the input builder (drop unpaired items; bypassed for store-mode
previous_response_id chaining where calls are matched server-side).
- Hard pre-send context gate: shrink oversized tool results and refuse
payloads that cannot fit the model's window before dispatch, so a
hopeless request is never sent or billed.
- Cap a single tool result entering the LLM context
(TOOL_RESULT_MAX_TOKENS, default 20000); the tool journal and
persistence keep the full result. Applied on the resume/continuation
path too.
- Skip the fallback attempt when the payload cannot fit the fallback
model's context window (10% estimation slack).
- Compression: never save a compression point that does not reduce
tokens; bound oversized verbatim fields kept after a compression
point (COMPRESSION_RECENT_FIELD_MAX_TOKENS, default 8000).
Robustness fixes from review:
- Google parallel function calls: complete index-less ToolCalls are no
longer merged into one another (dict arguments raised TypeError on
+=; second call could execute with the first call's arguments).
- Trailing-frame failures after a delivered answer no longer error the
stream or restream the whole answer from the fallback
(_stream_reached_finish).
- In-memory compression falls back to minimal pruning when the summary
is not smaller than the original.
- keep<=0 guard in the middle-truncation helpers (a tiny cap returned
marker + full text).
Frontend: tooltip on the Analytics tokens stat card explaining that
agent tool loops re-send conversation context on every step.
The classic prompt renderer now passes artifact_parent={conversation_id} so a
normal agent's prompt can resolve a prior-turn artifact with
{{ artifacts.artifact(id) }}, scoped to its own conversation (parent-derived
authz; a missing conversation_id safely yields an empty lookup). This makes the
artifacts variable surfaced in the agent builder functional, and is the parent
wiring an agent-as-tool / subagent feature would also rely on.
Let workflow runs consume and produce documents end to end: bridge uploaded
attachments into run-scoped artifacts so nodes receive the input documents
(with a per-run cap and server-computed size/sha256, and the run row pre-created
so produced artifacts are authorized during the run); emit the run id to the
client and add a builder panel that lists, previews, and downloads a run's
artifacts; and allow attaching documents to a Preview run via the existing
upload flow.
Also fixes issues a compliance workflow surfaced: attachment ownership now keys
on the raw identity instead of a sanitized one (the sanitized form could not be
read back and could collide across users); workflow code nodes read prior state
from a state.json data file instead of templating it into the program, so
untrusted document content can never be interpolated into executed code;
structured node output wrapped in code fences is recovered; and the live
speech-to-text ownership check compares the raw identity.
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.
Rewrite the default/creative/strict presets (classic + agentic) into
structured sections: grounding and cite-by-title guidance, insufficient-
context behavior, current date, respond-in-user-language, scoped mermaid
usage, an untrusted-content guardrail, and a conditional XML-tagged
document context block. A memory directory listing is injected at render
time via the template prefetch mechanism so the model starts oriented
without burning a tool call.
Fixes along the way:
- Agentic preset swap was dead code: _get_prompt_content cached the
classic preset before create_agent's swap check ran, so agentic and
research agents always got the classic preset. The swap now happens
inside _get_prompt_content.
- Jinja autoescape corrupted document content in custom prompts
(< -> <); prompts are not HTML, autoescape is now off.
- Literal {summaries} leaked into the prompt when no docs were
retrieved; the placeholder is now stripped.
- Agentic/research prompts referenced tool names from a dropped naming
scheme (search_internal, reason_think); they now reference the real
names (search, reason).
- The strict preset told the model to "be very creative and use your
imagination" right after "never make up information".
- extract_tool_usages recorded intermediate attribute chains as
bare-tool usages, which meant "run all actions" at prefetch; only
maximal chains are recorded now.
- Headless runs retrieved docs but never rendered them into the
prompt; the prompt is now rendered like the streaming path.
- Default tools were unreachable by name in prompt templates
(prefetch results were keyed by synthetic id only); defaults now
claim the name key unless an explicit row shadows it.
Tool layer: memory/notes/todo actions are namespaced (memory_view,
note_overwrite, todo_create, ...) with legacy unprefixed names still
accepted via prefix stripping; duplicate action names across tools are
disambiguated with the owning tool's name instead of numeric suffixes;
thin tool descriptions rewritten (brave, duckduckgo, telegram, ntfy,
cryptoprice, read_webpage, internal_search, think).
Docs are now wrapped per chunk in <document index>/<source>/<content>
tags for citation-by-title support.
* 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