diff --git a/server/internal/api/web/chat.go b/server/internal/api/web/chat.go --- a/server/internal/api/web/chat.go +++ b/server/internal/api/web/chat.go @@ -15,6 +15,7 @@ apimw "github.com/taciturnaxolotl/potluck/internal/api/middleware" "github.com/taciturnaxolotl/potluck/internal/pool" "github.com/taciturnaxolotl/potluck/internal/provider" "github.com/taciturnaxolotl/potluck/internal/store" + "github.com/taciturnaxolotl/potluck/internal/stream" "github.com/taciturnaxolotl/potluck/internal/tools" ) @@ -263,6 +264,7 @@ seq := int64(0) clientGone := false ctxDone := r.Context().Done() genCtx := context.Background() + bus := s.Hub.Subscriber(streamID) emit := func(event string, payload map[string]any) { seq++ @@ -275,6 +277,7 @@ Event: event, Data: string(b), CreatedAt: time.Now().Unix(), }) + bus.Publish(stream.Event{Seq: seq, Type: event, Raw: json.RawMessage(b)}) if clientGone { return } diff --git a/server/internal/api/web/web.go b/server/internal/api/web/web.go --- a/server/internal/api/web/web.go +++ b/server/internal/api/web/web.go @@ -503,7 +503,10 @@ w.WriteHeader(200) flusher, _ := w.(http.Flusher) emit := func(ev stream.Event) { - b, _ := json.Marshal(ev) + b := []byte(ev.Raw) + if len(b) == 0 { + b, _ = json.Marshal(ev) + } _, _ = w.Write([]byte("data: ")) _, _ = w.Write(b) _, _ = w.Write([]byte("\n\n")) @@ -512,7 +515,12 @@ flusher.Flush() } } - // 1) Replay durable chunks after the requested seq. + // 1) Subscribe to the live bus FIRST so we don't miss events between + // replay and subscribe. + bus := s.Hub.Subscriber(streamID) + ch, done, doneEv := bus.Subscribe(64) + + // 2) Replay durable chunks from DB (missed while disconnected). events, err := stream.Replay(r.Context(), s.Q, streamID, afterSeq) if err == nil { for _, ev := range events { @@ -523,10 +531,27 @@ return } } } + // Replay done; now drain any events the bus buffered since subscribe, + // then continue tailing live. + for { + select { + case ev, ok := <-ch: + if !ok { + return + } + if ev.Seq <= afterSeq { + continue + } + emit(ev) + afterSeq = ev.Seq + if ev.Type == "done" || ev.Type == "error" { + return + } + default: + } + break + } - // 2) Attach to live bus for tailing events. - bus := s.Hub.Subscriber(streamID) - ch, done, doneEv := bus.Subscribe(64) if done { if doneEv.Seq > afterSeq { emit(doneEv) diff --git a/server/internal/stream/tee.go b/server/internal/stream/tee.go --- a/server/internal/stream/tee.go +++ b/server/internal/stream/tee.go @@ -194,6 +194,7 @@ var ev Event _ = json.Unmarshal([]byte(r.Data), &ev) ev.Seq = r.Seq ev.Type = r.Event + ev.Raw = json.RawMessage(r.Data) out = append(out, ev) } return out, nil diff --git a/web/src/routes/chat/+page.svelte b/web/src/routes/chat/+page.svelte --- a/web/src/routes/chat/+page.svelte +++ b/web/src/routes/chat/+page.svelte @@ -186,6 +186,8 @@ streaming = true; activeStreamId = resumeStreamId; streamingMsgId = resumeAssistantId; activeAfterSeq = Number.isFinite(resumeAfterSeq) ? resumeAfterSeq : 0; + streamStartMs = Date.now(); + streamFirstTokenMs = 0; attachStreamConsumer(resumeStreamId, resumeConvId, resumeAssistantId, Math.floor(Date.now() / 1000), resumeAfterSeq); } }); @@ -430,7 +432,6 @@ let resolvedUserId = tmpUserId; let resolvedAssistantId = tmpAssistantId; let accContent = ''; - let handedOffToStreamConsumer = false; try { const res = await fetch('/api/chat', { @@ -475,6 +476,14 @@ try { ev = JSON.parse(data); } catch { continue; + } + + // Keep activeAfterSeq current so resume picks up from the right spot. + if (typeof ev.seq === 'number' && ev.seq > activeAfterSeq) { + activeAfterSeq = ev.seq; + if (activeStreamId) { + localStorage.setItem('chat:active_stream_after_seq', String(activeAfterSeq)); + } } switch (ev.type) { @@ -534,15 +543,13 @@ activeConvId = serverConvId; goto(`/chat?c=${serverConvId}`, { replaceState: true, noScroll: true, keepFocus: true }); } + // Save stream metadata for resume-on-reload — don't hand off mid-stream. if (serverStreamId) { activeStreamId = serverStreamId; localStorage.setItem('chat:active_stream_id', serverStreamId); localStorage.setItem('chat:active_stream_conv_id', convId); localStorage.setItem('chat:active_stream_assistant_id', resolvedAssistantId); localStorage.setItem('chat:active_stream_after_seq', String(activeAfterSeq)); - attachStreamConsumer(serverStreamId, convId, resolvedAssistantId, now, activeAfterSeq); - handedOffToStreamConsumer = true; - break outer; } break; @@ -640,10 +647,8 @@ } catch (err) { await db.messages.delete(resolvedAssistantId); errorMsg = err instanceof Error ? err.message : 'Something went wrong'; } finally { - if (!handedOffToStreamConsumer) { - streaming = false; - streamingMsgId = null; - } + streaming = false; + streamingMsgId = null; } }