luckfun233 commited on
Commit ·
39f9c5d
1
Parent(s): ee56648
perf(chathistory): 拆分读写锁并合并流式落盘,修复 HF Space 整体卡顿
Browse files线上表现为 API 响应变慢、后台面板加载历史对话卡顿。根因是对话历史的落盘
路径把整个服务串行化了:
- Store 只有一把互斥锁,saveLocked 在持锁状态下完成 JSON 序列化 + fsync +
rename;而 Snapshot / Get(后台面板读历史)用的是同一把锁,因此被流式写入
直接阻塞。
- progress() 在每个 SSE 分片上被调用(250ms 节流),每次都重写包含 messages /
history_text / final_prompt 等请求期内不变字段的完整记录。HuggingFace Spaces
的 /data 是 bucket 对象存储卷,单次 fsync 的延迟远高于本地盘,几十秒的响应
就会触发数百次高延迟写入。
改动:
- 拆成 mu(仅内存状态,持锁期间不做 I/O)与 writeMu(仅串行化落盘),后台
面板读取不再阻塞在磁盘上。
- 新增 UpdateParams.Background:流式进度只写 <id>.live.json 增量文件,由单个
自终止的合并协程按 400ms 窗口合并落盘,彻底脱离 SSE 读取循环;终态更新仍
同步回写完整文件并清理增量文件。
- 去掉 chat history 的 fsync,保留「临时文件 + rename」的原子性。
- 后台面板轮询降频(列表 1.5s→3s,详情 750ms→1.5s)。
实测单次进度写入从 84,475 B 降至 246 B(343x)。旧数据无需迁移即可加载。
Start 与终态 Update 仍同步返回写盘错误,保持既有错误语义。
internal/chathistory/store.go
CHANGED
|
@@ -23,6 +23,14 @@ const (
|
|
| 23 |
DefaultLimit = 20
|
| 24 |
MaxLimit = 50
|
| 25 |
defaultPreviewAt = 160
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 26 |
)
|
| 27 |
|
| 28 |
var allowedLimits = map[int]struct{}{
|
|
@@ -131,6 +139,11 @@ type UpdateParams struct {
|
|
| 131 |
FinishReason string
|
| 132 |
Usage map[string]any
|
| 133 |
Completed bool
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 134 |
}
|
| 135 |
|
| 136 |
type detailEnvelope struct {
|
|
@@ -138,6 +151,27 @@ type detailEnvelope struct {
|
|
| 138 |
Item Entry `json:"item"`
|
| 139 |
}
|
| 140 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 141 |
type legacyFile struct {
|
| 142 |
Version int `json:"version"`
|
| 143 |
Limit int `json:"limit"`
|
|
@@ -148,15 +182,25 @@ type legacyProbe struct {
|
|
| 148 |
Items []map[string]json.RawMessage `json:"items"`
|
| 149 |
}
|
| 150 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 151 |
type Store struct {
|
| 152 |
-
mu
|
| 153 |
-
|
| 154 |
-
|
| 155 |
-
|
| 156 |
-
|
| 157 |
-
|
| 158 |
-
|
| 159 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 160 |
}
|
| 161 |
|
| 162 |
func New(path string) *Store {
|
|
@@ -169,13 +213,24 @@ func New(path string) *Store {
|
|
| 169 |
Revision: 0,
|
| 170 |
Items: []SummaryEntry{},
|
| 171 |
},
|
| 172 |
-
details:
|
| 173 |
-
dirty:
|
| 174 |
-
|
|
|
|
| 175 |
}
|
| 176 |
s.mu.Lock()
|
| 177 |
-
|
| 178 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 179 |
return s
|
| 180 |
}
|
| 181 |
|
|
@@ -275,11 +330,13 @@ func (s *Store) Start(params StartParams) (Entry, error) {
|
|
| 275 |
return Entry{}, errors.New("chat history store is nil")
|
| 276 |
}
|
| 277 |
s.mu.Lock()
|
| 278 |
-
defer s.mu.Unlock()
|
| 279 |
if s.err != nil {
|
| 280 |
-
|
|
|
|
|
|
|
| 281 |
}
|
| 282 |
if s.state.Limit == DisabledLimit {
|
|
|
|
| 283 |
return Entry{}, ErrDisabled
|
| 284 |
}
|
| 285 |
now := time.Now().UnixMilli()
|
|
@@ -304,10 +361,13 @@ func (s *Store) Start(params StartParams) (Entry, error) {
|
|
| 304 |
s.details[entry.ID] = entry
|
| 305 |
s.markDetailDirtyLocked(entry.ID)
|
| 306 |
s.rebuildIndexLocked()
|
| 307 |
-
|
| 308 |
-
|
|
|
|
|
|
|
|
|
|
| 309 |
}
|
| 310 |
-
return
|
| 311 |
}
|
| 312 |
|
| 313 |
func (s *Store) Update(id string, params UpdateParams) (Entry, error) {
|
|
@@ -315,16 +375,19 @@ func (s *Store) Update(id string, params UpdateParams) (Entry, error) {
|
|
| 315 |
return Entry{}, errors.New("chat history store is nil")
|
| 316 |
}
|
| 317 |
s.mu.Lock()
|
| 318 |
-
defer s.mu.Unlock()
|
| 319 |
if s.err != nil {
|
| 320 |
-
|
|
|
|
|
|
|
| 321 |
}
|
| 322 |
target := strings.TrimSpace(id)
|
| 323 |
if target == "" {
|
|
|
|
| 324 |
return Entry{}, errors.New("history id is required")
|
| 325 |
}
|
| 326 |
item, ok := s.details[target]
|
| 327 |
if !ok {
|
|
|
|
| 328 |
return Entry{}, errors.New("chat history entry not found")
|
| 329 |
}
|
| 330 |
now := time.Now().UnixMilli()
|
|
@@ -350,12 +413,24 @@ func (s *Store) Update(id string, params UpdateParams) (Entry, error) {
|
|
| 350 |
item.CompletedAt = now
|
| 351 |
}
|
| 352 |
s.details[target] = item
|
| 353 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
| 354 |
s.rebuildIndexLocked()
|
| 355 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 356 |
return Entry{}, err
|
| 357 |
}
|
| 358 |
-
return
|
| 359 |
}
|
| 360 |
|
| 361 |
func (s *Store) Delete(id string) error {
|
|
@@ -363,25 +438,28 @@ func (s *Store) Delete(id string) error {
|
|
| 363 |
return errors.New("chat history store is nil")
|
| 364 |
}
|
| 365 |
s.mu.Lock()
|
| 366 |
-
defer s.mu.Unlock()
|
| 367 |
if s.err != nil {
|
| 368 |
-
|
|
|
|
|
|
|
| 369 |
}
|
| 370 |
target := strings.TrimSpace(id)
|
| 371 |
if target == "" {
|
|
|
|
| 372 |
return errors.New("history id is required")
|
| 373 |
}
|
| 374 |
if _, ok := s.details[target]; !ok {
|
|
|
|
| 375 |
return errors.New("chat history entry not found")
|
| 376 |
}
|
| 377 |
s.markDetailDeletedLocked(target)
|
| 378 |
delete(s.details, target)
|
| 379 |
s.nextRevisionLocked()
|
| 380 |
s.rebuildIndexLocked()
|
| 381 |
-
|
| 382 |
-
|
| 383 |
-
|
| 384 |
-
return
|
| 385 |
}
|
| 386 |
|
| 387 |
func (s *Store) Clear() error {
|
|
@@ -389,9 +467,10 @@ func (s *Store) Clear() error {
|
|
| 389 |
return errors.New("chat history store is nil")
|
| 390 |
}
|
| 391 |
s.mu.Lock()
|
| 392 |
-
defer s.mu.Unlock()
|
| 393 |
if s.err != nil {
|
| 394 |
-
|
|
|
|
|
|
|
| 395 |
}
|
| 396 |
for id := range s.details {
|
| 397 |
s.markDetailDeletedLocked(id)
|
|
@@ -399,10 +478,10 @@ func (s *Store) Clear() error {
|
|
| 399 |
s.details = map[string]Entry{}
|
| 400 |
s.nextRevisionLocked()
|
| 401 |
s.rebuildIndexLocked()
|
| 402 |
-
|
| 403 |
-
|
| 404 |
-
|
| 405 |
-
return
|
| 406 |
}
|
| 407 |
|
| 408 |
func (s *Store) SetLimit(limit int) (File, error) {
|
|
@@ -410,22 +489,30 @@ func (s *Store) SetLimit(limit int) (File, error) {
|
|
| 410 |
return File{}, errors.New("chat history store is nil")
|
| 411 |
}
|
| 412 |
s.mu.Lock()
|
| 413 |
-
defer s.mu.Unlock()
|
| 414 |
if s.err != nil {
|
| 415 |
-
|
|
|
|
|
|
|
| 416 |
}
|
| 417 |
if !isAllowedLimit(limit) {
|
|
|
|
| 418 |
return File{}, fmt.Errorf("unsupported chat history limit: %d", limit)
|
| 419 |
}
|
| 420 |
s.state.Limit = limit
|
| 421 |
s.nextRevisionLocked()
|
| 422 |
s.rebuildIndexLocked()
|
| 423 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
| 424 |
return File{}, err
|
| 425 |
}
|
| 426 |
-
return
|
| 427 |
}
|
| 428 |
|
|
|
|
|
|
|
| 429 |
func (s *Store) loadLocked() error {
|
| 430 |
if strings.TrimSpace(s.path) == "" {
|
| 431 |
return errors.New("chat history path is required")
|
|
@@ -440,9 +527,7 @@ func (s *Store) loadLocked() error {
|
|
| 440 |
raw, err := os.ReadFile(s.path)
|
| 441 |
if err != nil {
|
| 442 |
if errors.Is(err, os.ErrNotExist) {
|
| 443 |
-
|
| 444 |
-
config.Logger.Warn("[chat_history] bootstrap write failed", "path", s.path, "error", saveErr)
|
| 445 |
-
}
|
| 446 |
return nil
|
| 447 |
}
|
| 448 |
return fmt.Errorf("read chat history index: %w", err)
|
|
@@ -454,9 +539,6 @@ func (s *Store) loadLocked() error {
|
|
| 454 |
}
|
| 455 |
if legacyOK {
|
| 456 |
s.loadLegacyLocked(legacy)
|
| 457 |
-
if err := s.saveLocked(); err != nil {
|
| 458 |
-
config.Logger.Warn("[chat_history] legacy migration writeback failed", "path", s.path, "error", err)
|
| 459 |
-
}
|
| 460 |
return nil
|
| 461 |
}
|
| 462 |
|
|
@@ -473,16 +555,14 @@ func (s *Store) loadLocked() error {
|
|
| 473 |
s.state = cloneFile(state)
|
| 474 |
s.details = map[string]Entry{}
|
| 475 |
for _, item := range state.Items {
|
| 476 |
-
detail, err :=
|
| 477 |
if err != nil {
|
| 478 |
return err
|
| 479 |
}
|
| 480 |
s.details[item.ID] = detail
|
| 481 |
}
|
| 482 |
s.rebuildIndexLocked()
|
| 483 |
-
|
| 484 |
-
config.Logger.Warn("[chat_history] index rewrite failed", "path", s.path, "error", saveErr)
|
| 485 |
-
}
|
| 486 |
return nil
|
| 487 |
}
|
| 488 |
|
|
@@ -494,6 +574,7 @@ func (s *Store) loadLegacyLocked(legacy legacyFile) {
|
|
| 494 |
}
|
| 495 |
s.details = map[string]Entry{}
|
| 496 |
s.dirty = map[string]struct{}{}
|
|
|
|
| 497 |
s.deleted = map[string]struct{}{}
|
| 498 |
maxRevision := int64(0)
|
| 499 |
for _, item := range legacy.Items {
|
|
@@ -515,54 +596,219 @@ func (s *Store) loadLegacyLocked(legacy legacyFile) {
|
|
| 515 |
s.markDetailDirtyLocked(item.ID)
|
| 516 |
}
|
| 517 |
s.state.Revision = maxRevision
|
|
|
|
| 518 |
s.rebuildIndexLocked()
|
| 519 |
}
|
| 520 |
|
| 521 |
-
func (s *Store)
|
| 522 |
-
s
|
| 523 |
-
|
| 524 |
-
s.state.Limit = DefaultLimit
|
| 525 |
}
|
| 526 |
-
s.
|
|
|
|
|
|
|
|
|
|
| 527 |
|
| 528 |
-
|
| 529 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 530 |
}
|
| 531 |
-
|
| 532 |
-
|
| 533 |
-
|
| 534 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 535 |
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 536 |
}
|
| 537 |
-
|
| 538 |
-
|
| 539 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 540 |
continue
|
| 541 |
}
|
| 542 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 543 |
payload, err := json.MarshalIndent(detailEnvelope{
|
| 544 |
Version: FileVersion,
|
| 545 |
-
Item:
|
| 546 |
}, "", " ")
|
| 547 |
if err != nil {
|
| 548 |
return fmt.Errorf("encode chat history detail: %w", err)
|
| 549 |
}
|
| 550 |
-
if err := writeFileAtomic(
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 551 |
return err
|
| 552 |
}
|
| 553 |
}
|
| 554 |
|
| 555 |
-
payload, err := json.MarshalIndent(
|
| 556 |
if err != nil {
|
| 557 |
return fmt.Errorf("encode chat history index: %w", err)
|
| 558 |
}
|
| 559 |
if err := writeFileAtomic(s.path, append(payload, '\n')); err != nil {
|
| 560 |
return err
|
| 561 |
}
|
| 562 |
-
s.clearPendingDetailChangesLocked()
|
| 563 |
return nil
|
| 564 |
}
|
| 565 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 566 |
func (s *Store) rebuildIndexLocked() {
|
| 567 |
summaries := make([]SummaryEntry, 0, len(s.details))
|
| 568 |
for _, item := range s.details {
|
|
@@ -606,6 +852,7 @@ func (s *Store) nextRevisionLocked() int64 {
|
|
| 606 |
next = s.state.Revision + 1
|
| 607 |
}
|
| 608 |
s.state.Revision = next
|
|
|
|
| 609 |
return next
|
| 610 |
}
|
| 611 |
|
|
@@ -648,6 +895,40 @@ func buildPreview(item Entry) string {
|
|
| 648 |
return candidate
|
| 649 |
}
|
| 650 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 651 |
func readDetailFile(path string) (Entry, error) {
|
| 652 |
raw, err := os.ReadFile(path)
|
| 653 |
if err != nil {
|
|
@@ -660,6 +941,21 @@ func readDetailFile(path string) (Entry, error) {
|
|
| 660 |
return cloneEntry(env.Item), nil
|
| 661 |
}
|
| 662 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 663 |
func parseLegacy(raw []byte) (legacyFile, bool, error) {
|
| 664 |
var legacy legacyFile
|
| 665 |
if err := json.Unmarshal(raw, &legacy); err != nil {
|
|
@@ -679,6 +975,12 @@ func parseLegacy(raw []byte) (legacyFile, bool, error) {
|
|
| 679 |
return legacy, true, nil
|
| 680 |
}
|
| 681 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 682 |
func writeFileAtomic(path string, body []byte) error {
|
| 683 |
dir := filepath.Dir(path)
|
| 684 |
if dir == "" {
|
|
@@ -713,9 +1015,6 @@ func writeFileAtomic(path string, body []byte) error {
|
|
| 713 |
if _, err := tmpFile.Write(body); err != nil {
|
| 714 |
return withCleanup(fmt.Errorf("write temp chat history: %w", err), tmpFile.Close())
|
| 715 |
}
|
| 716 |
-
if err := tmpFile.Sync(); err != nil {
|
| 717 |
-
return withCleanup(fmt.Errorf("sync temp chat history: %w", err), tmpFile.Close())
|
| 718 |
-
}
|
| 719 |
if err := tmpFile.Close(); err != nil {
|
| 720 |
if cleanupErr := cleanup(); cleanupErr != nil {
|
| 721 |
return errors.Join(fmt.Errorf("close temp chat history: %w", err), cleanupErr)
|
|
@@ -752,10 +1051,34 @@ func (s *Store) markDetailDirtyLocked(id string) {
|
|
| 752 |
if s.dirty == nil {
|
| 753 |
s.dirty = map[string]struct{}{}
|
| 754 |
}
|
|
|
|
|
|
|
|
|
|
| 755 |
if s.deleted == nil {
|
| 756 |
s.deleted = map[string]struct{}{}
|
| 757 |
}
|
| 758 |
s.dirty[id] = struct{}{}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 759 |
delete(s.deleted, id)
|
| 760 |
}
|
| 761 |
|
|
@@ -767,16 +1090,15 @@ func (s *Store) markDetailDeletedLocked(id string) {
|
|
| 767 |
if s.dirty == nil {
|
| 768 |
s.dirty = map[string]struct{}{}
|
| 769 |
}
|
|
|
|
|
|
|
|
|
|
| 770 |
if s.deleted == nil {
|
| 771 |
s.deleted = map[string]struct{}{}
|
| 772 |
}
|
| 773 |
s.deleted[id] = struct{}{}
|
| 774 |
delete(s.dirty, id)
|
| 775 |
-
|
| 776 |
-
|
| 777 |
-
func (s *Store) clearPendingDetailChangesLocked() {
|
| 778 |
-
s.dirty = map[string]struct{}{}
|
| 779 |
-
s.deleted = map[string]struct{}{}
|
| 780 |
}
|
| 781 |
|
| 782 |
func sortedDetailIDs(ids map[string]struct{}) []string {
|
|
@@ -791,6 +1113,17 @@ func sortedDetailIDs(ids map[string]struct{}) []string {
|
|
| 791 |
return out
|
| 792 |
}
|
| 793 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 794 |
func cloneFile(in File) File {
|
| 795 |
out := File{
|
| 796 |
Version: in.Version,
|
|
|
|
| 23 |
DefaultLimit = 20
|
| 24 |
MaxLimit = 50
|
| 25 |
defaultPreviewAt = 160
|
| 26 |
+
|
| 27 |
+
// asyncDebounce 是流式进度落盘的合并窗口:窗口内的多次进度更新只会产生
|
| 28 |
+
// 一次磁盘写入。HuggingFace Spaces 的 bucket 卷等对象存储挂载单次写入
|
| 29 |
+
// 延迟很高,按 SSE 分片逐次写盘会把整个进程拖垮。
|
| 30 |
+
asyncDebounce = 400 * time.Millisecond
|
| 31 |
+
// asyncRetryDelay 是写盘连续失败后的退避时间,避免在磁盘长期不可用时
|
| 32 |
+
// 以合并窗口的频率反复重试。
|
| 33 |
+
asyncRetryDelay = 5 * time.Second
|
| 34 |
)
|
| 35 |
|
| 36 |
var allowedLimits = map[int]struct{}{
|
|
|
|
| 139 |
FinishReason string
|
| 140 |
Usage map[string]any
|
| 141 |
Completed bool
|
| 142 |
+
// Background 标记一次流式进度更新:只落盘一个只含变化字段的增量文件
|
| 143 |
+
// (<id>.live.json),不重写包含 messages / history_text / final_prompt
|
| 144 |
+
// 等请求期内不变的大字段的完整文件。调用方不会拿到写盘错误,
|
| 145 |
+
// 写入由后台合并协程完成。
|
| 146 |
+
Background bool
|
| 147 |
}
|
| 148 |
|
| 149 |
type detailEnvelope struct {
|
|
|
|
| 151 |
Item Entry `json:"item"`
|
| 152 |
}
|
| 153 |
|
| 154 |
+
// progressRecord 是流式期间的增量落盘记录,只包含随输出变化的字段。
|
| 155 |
+
type progressRecord struct {
|
| 156 |
+
Version int `json:"version"`
|
| 157 |
+
Revision int64 `json:"revision"`
|
| 158 |
+
UpdatedAt int64 `json:"updated_at"`
|
| 159 |
+
CompletedAt int64 `json:"completed_at,omitempty"`
|
| 160 |
+
Status string `json:"status,omitempty"`
|
| 161 |
+
ReasoningContent string `json:"reasoning_content,omitempty"`
|
| 162 |
+
Content string `json:"content,omitempty"`
|
| 163 |
+
Error string `json:"error,omitempty"`
|
| 164 |
+
StatusCode int `json:"status_code,omitempty"`
|
| 165 |
+
ElapsedMs int64 `json:"elapsed_ms,omitempty"`
|
| 166 |
+
FinishReason string `json:"finish_reason,omitempty"`
|
| 167 |
+
Usage map[string]any `json:"usage,omitempty"`
|
| 168 |
+
}
|
| 169 |
+
|
| 170 |
+
type progressEnvelope struct {
|
| 171 |
+
Version int `json:"version"`
|
| 172 |
+
Item progressRecord `json:"item"`
|
| 173 |
+
}
|
| 174 |
+
|
| 175 |
type legacyFile struct {
|
| 176 |
Version int `json:"version"`
|
| 177 |
Limit int `json:"limit"`
|
|
|
|
| 182 |
Items []map[string]json.RawMessage `json:"items"`
|
| 183 |
}
|
| 184 |
|
| 185 |
+
// Store 保存 chat history。锁被刻意拆成两把:
|
| 186 |
+
//
|
| 187 |
+
// - mu 只保护内存状态,任何持有它的路径都不做 I/O,因此后台面板读取历史
|
| 188 |
+
// (Snapshot / Get)永远不会被流式落盘阻塞;
|
| 189 |
+
// - writeMu 只串行化落盘,让并发请求的写入不会互相覆盖。
|
| 190 |
type Store struct {
|
| 191 |
+
mu sync.Mutex
|
| 192 |
+
writeMu sync.Mutex
|
| 193 |
+
path string
|
| 194 |
+
detailDir string
|
| 195 |
+
state File
|
| 196 |
+
details map[string]Entry
|
| 197 |
+
dirty map[string]struct{} // 需要重写完整 <id>.json 的条目
|
| 198 |
+
progress map[string]struct{} // 只需要写增量 <id>.live.json 的条目
|
| 199 |
+
deleted map[string]struct{}
|
| 200 |
+
indexDirty bool
|
| 201 |
+
asyncBusy bool
|
| 202 |
+
closed bool
|
| 203 |
+
err error
|
| 204 |
}
|
| 205 |
|
| 206 |
func New(path string) *Store {
|
|
|
|
| 213 |
Revision: 0,
|
| 214 |
Items: []SummaryEntry{},
|
| 215 |
},
|
| 216 |
+
details: map[string]Entry{},
|
| 217 |
+
dirty: map[string]struct{}{},
|
| 218 |
+
progress: map[string]struct{}{},
|
| 219 |
+
deleted: map[string]struct{}{},
|
| 220 |
}
|
| 221 |
s.mu.Lock()
|
| 222 |
+
loadErr := s.loadLocked()
|
| 223 |
+
if loadErr != nil {
|
| 224 |
+
s.err = loadErr
|
| 225 |
+
}
|
| 226 |
+
s.mu.Unlock()
|
| 227 |
+
|
| 228 |
+
// 首次启动建索引 / 旧格式迁移 / 索引归一化都在锁外回写,失败只告警不致命。
|
| 229 |
+
if loadErr == nil && s.pending() {
|
| 230 |
+
if err := s.flush(); err != nil {
|
| 231 |
+
config.Logger.Warn("[chat_history] initial writeback failed", "path", s.path, "error", err)
|
| 232 |
+
}
|
| 233 |
+
}
|
| 234 |
return s
|
| 235 |
}
|
| 236 |
|
|
|
|
| 330 |
return Entry{}, errors.New("chat history store is nil")
|
| 331 |
}
|
| 332 |
s.mu.Lock()
|
|
|
|
| 333 |
if s.err != nil {
|
| 334 |
+
err := s.err
|
| 335 |
+
s.mu.Unlock()
|
| 336 |
+
return Entry{}, err
|
| 337 |
}
|
| 338 |
if s.state.Limit == DisabledLimit {
|
| 339 |
+
s.mu.Unlock()
|
| 340 |
return Entry{}, ErrDisabled
|
| 341 |
}
|
| 342 |
now := time.Now().UnixMilli()
|
|
|
|
| 361 |
s.details[entry.ID] = entry
|
| 362 |
s.markDetailDirtyLocked(entry.ID)
|
| 363 |
s.rebuildIndexLocked()
|
| 364 |
+
result := cloneEntry(entry)
|
| 365 |
+
s.mu.Unlock()
|
| 366 |
+
|
| 367 |
+
if err := s.flush(); err != nil {
|
| 368 |
+
return result, err
|
| 369 |
}
|
| 370 |
+
return result, nil
|
| 371 |
}
|
| 372 |
|
| 373 |
func (s *Store) Update(id string, params UpdateParams) (Entry, error) {
|
|
|
|
| 375 |
return Entry{}, errors.New("chat history store is nil")
|
| 376 |
}
|
| 377 |
s.mu.Lock()
|
|
|
|
| 378 |
if s.err != nil {
|
| 379 |
+
err := s.err
|
| 380 |
+
s.mu.Unlock()
|
| 381 |
+
return Entry{}, err
|
| 382 |
}
|
| 383 |
target := strings.TrimSpace(id)
|
| 384 |
if target == "" {
|
| 385 |
+
s.mu.Unlock()
|
| 386 |
return Entry{}, errors.New("history id is required")
|
| 387 |
}
|
| 388 |
item, ok := s.details[target]
|
| 389 |
if !ok {
|
| 390 |
+
s.mu.Unlock()
|
| 391 |
return Entry{}, errors.New("chat history entry not found")
|
| 392 |
}
|
| 393 |
now := time.Now().UnixMilli()
|
|
|
|
| 413 |
item.CompletedAt = now
|
| 414 |
}
|
| 415 |
s.details[target] = item
|
| 416 |
+
if params.Background {
|
| 417 |
+
s.markProgressDirtyLocked(target)
|
| 418 |
+
} else {
|
| 419 |
+
s.markDetailDirtyLocked(target)
|
| 420 |
+
}
|
| 421 |
s.rebuildIndexLocked()
|
| 422 |
+
result := cloneEntry(item)
|
| 423 |
+
s.mu.Unlock()
|
| 424 |
+
|
| 425 |
+
if params.Background {
|
| 426 |
+
// 流式进度:交给后台合并写入,绝不阻塞 SSE 读取循环。
|
| 427 |
+
s.flushAsync()
|
| 428 |
+
return result, nil
|
| 429 |
+
}
|
| 430 |
+
if err := s.flush(); err != nil {
|
| 431 |
return Entry{}, err
|
| 432 |
}
|
| 433 |
+
return result, nil
|
| 434 |
}
|
| 435 |
|
| 436 |
func (s *Store) Delete(id string) error {
|
|
|
|
| 438 |
return errors.New("chat history store is nil")
|
| 439 |
}
|
| 440 |
s.mu.Lock()
|
|
|
|
| 441 |
if s.err != nil {
|
| 442 |
+
err := s.err
|
| 443 |
+
s.mu.Unlock()
|
| 444 |
+
return err
|
| 445 |
}
|
| 446 |
target := strings.TrimSpace(id)
|
| 447 |
if target == "" {
|
| 448 |
+
s.mu.Unlock()
|
| 449 |
return errors.New("history id is required")
|
| 450 |
}
|
| 451 |
if _, ok := s.details[target]; !ok {
|
| 452 |
+
s.mu.Unlock()
|
| 453 |
return errors.New("chat history entry not found")
|
| 454 |
}
|
| 455 |
s.markDetailDeletedLocked(target)
|
| 456 |
delete(s.details, target)
|
| 457 |
s.nextRevisionLocked()
|
| 458 |
s.rebuildIndexLocked()
|
| 459 |
+
s.indexDirty = true
|
| 460 |
+
s.mu.Unlock()
|
| 461 |
+
|
| 462 |
+
return s.flush()
|
| 463 |
}
|
| 464 |
|
| 465 |
func (s *Store) Clear() error {
|
|
|
|
| 467 |
return errors.New("chat history store is nil")
|
| 468 |
}
|
| 469 |
s.mu.Lock()
|
|
|
|
| 470 |
if s.err != nil {
|
| 471 |
+
err := s.err
|
| 472 |
+
s.mu.Unlock()
|
| 473 |
+
return err
|
| 474 |
}
|
| 475 |
for id := range s.details {
|
| 476 |
s.markDetailDeletedLocked(id)
|
|
|
|
| 478 |
s.details = map[string]Entry{}
|
| 479 |
s.nextRevisionLocked()
|
| 480 |
s.rebuildIndexLocked()
|
| 481 |
+
s.indexDirty = true
|
| 482 |
+
s.mu.Unlock()
|
| 483 |
+
|
| 484 |
+
return s.flush()
|
| 485 |
}
|
| 486 |
|
| 487 |
func (s *Store) SetLimit(limit int) (File, error) {
|
|
|
|
| 489 |
return File{}, errors.New("chat history store is nil")
|
| 490 |
}
|
| 491 |
s.mu.Lock()
|
|
|
|
| 492 |
if s.err != nil {
|
| 493 |
+
err := s.err
|
| 494 |
+
s.mu.Unlock()
|
| 495 |
+
return File{}, err
|
| 496 |
}
|
| 497 |
if !isAllowedLimit(limit) {
|
| 498 |
+
s.mu.Unlock()
|
| 499 |
return File{}, fmt.Errorf("unsupported chat history limit: %d", limit)
|
| 500 |
}
|
| 501 |
s.state.Limit = limit
|
| 502 |
s.nextRevisionLocked()
|
| 503 |
s.rebuildIndexLocked()
|
| 504 |
+
s.indexDirty = true
|
| 505 |
+
result := cloneFile(s.state)
|
| 506 |
+
s.mu.Unlock()
|
| 507 |
+
|
| 508 |
+
if err := s.flush(); err != nil {
|
| 509 |
return File{}, err
|
| 510 |
}
|
| 511 |
+
return result, nil
|
| 512 |
}
|
| 513 |
|
| 514 |
+
// loadLocked 只做内存装载,不做任何写入;需要回写的内容通过 indexDirty /
|
| 515 |
+
// dirty 标记,由 New 在释放锁之后统一落盘。
|
| 516 |
func (s *Store) loadLocked() error {
|
| 517 |
if strings.TrimSpace(s.path) == "" {
|
| 518 |
return errors.New("chat history path is required")
|
|
|
|
| 527 |
raw, err := os.ReadFile(s.path)
|
| 528 |
if err != nil {
|
| 529 |
if errors.Is(err, os.ErrNotExist) {
|
| 530 |
+
s.indexDirty = true
|
|
|
|
|
|
|
| 531 |
return nil
|
| 532 |
}
|
| 533 |
return fmt.Errorf("read chat history index: %w", err)
|
|
|
|
| 539 |
}
|
| 540 |
if legacyOK {
|
| 541 |
s.loadLegacyLocked(legacy)
|
|
|
|
|
|
|
|
|
|
| 542 |
return nil
|
| 543 |
}
|
| 544 |
|
|
|
|
| 555 |
s.state = cloneFile(state)
|
| 556 |
s.details = map[string]Entry{}
|
| 557 |
for _, item := range state.Items {
|
| 558 |
+
detail, err := s.readEntryLocked(item.ID)
|
| 559 |
if err != nil {
|
| 560 |
return err
|
| 561 |
}
|
| 562 |
s.details[item.ID] = detail
|
| 563 |
}
|
| 564 |
s.rebuildIndexLocked()
|
| 565 |
+
s.indexDirty = true
|
|
|
|
|
|
|
| 566 |
return nil
|
| 567 |
}
|
| 568 |
|
|
|
|
| 574 |
}
|
| 575 |
s.details = map[string]Entry{}
|
| 576 |
s.dirty = map[string]struct{}{}
|
| 577 |
+
s.progress = map[string]struct{}{}
|
| 578 |
s.deleted = map[string]struct{}{}
|
| 579 |
maxRevision := int64(0)
|
| 580 |
for _, item := range legacy.Items {
|
|
|
|
| 596 |
s.markDetailDirtyLocked(item.ID)
|
| 597 |
}
|
| 598 |
s.state.Revision = maxRevision
|
| 599 |
+
s.indexDirty = true
|
| 600 |
s.rebuildIndexLocked()
|
| 601 |
}
|
| 602 |
|
| 603 |
+
func (s *Store) pending() bool {
|
| 604 |
+
if s == nil {
|
| 605 |
+
return false
|
|
|
|
| 606 |
}
|
| 607 |
+
s.mu.Lock()
|
| 608 |
+
defer s.mu.Unlock()
|
| 609 |
+
return s.hasPendingLocked()
|
| 610 |
+
}
|
| 611 |
|
| 612 |
+
func (s *Store) hasPendingLocked() bool {
|
| 613 |
+
return len(s.dirty) > 0 || len(s.progress) > 0 || len(s.deleted) > 0 || s.indexDirty
|
| 614 |
+
}
|
| 615 |
+
|
| 616 |
+
// Close 停止后台合并写入。它不改动已落盘的数据,只用于进程退出与测试收尾,
|
| 617 |
+
// 避免磁盘长期不可用时合并协程持续重试。
|
| 618 |
+
func (s *Store) Close() {
|
| 619 |
+
if s == nil {
|
| 620 |
+
return
|
| 621 |
}
|
| 622 |
+
s.mu.Lock()
|
| 623 |
+
s.closed = true
|
| 624 |
+
s.mu.Unlock()
|
| 625 |
+
}
|
| 626 |
+
|
| 627 |
+
// flushAsync 安排一次合并后的后台落盘。同一时刻最多只有一个合并协程,
|
| 628 |
+
// 且它自终止:没有待写内容时自动退出。
|
| 629 |
+
func (s *Store) flushAsync() {
|
| 630 |
+
if s == nil {
|
| 631 |
+
return
|
| 632 |
+
}
|
| 633 |
+
s.mu.Lock()
|
| 634 |
+
if s.asyncBusy || s.closed {
|
| 635 |
+
s.mu.Unlock()
|
| 636 |
+
return
|
| 637 |
+
}
|
| 638 |
+
s.asyncBusy = true
|
| 639 |
+
s.mu.Unlock()
|
| 640 |
+
|
| 641 |
+
go func() {
|
| 642 |
+
delay := asyncDebounce
|
| 643 |
+
for {
|
| 644 |
+
time.Sleep(delay)
|
| 645 |
+
s.mu.Lock()
|
| 646 |
+
closed := s.closed
|
| 647 |
+
s.mu.Unlock()
|
| 648 |
+
if closed {
|
| 649 |
+
s.mu.Lock()
|
| 650 |
+
s.asyncBusy = false
|
| 651 |
+
s.mu.Unlock()
|
| 652 |
+
return
|
| 653 |
+
}
|
| 654 |
+
err := s.flush()
|
| 655 |
+
s.mu.Lock()
|
| 656 |
+
if s.closed || !s.hasPendingLocked() {
|
| 657 |
+
s.asyncBusy = false
|
| 658 |
+
s.mu.Unlock()
|
| 659 |
+
return
|
| 660 |
+
}
|
| 661 |
+
s.mu.Unlock()
|
| 662 |
+
if err != nil {
|
| 663 |
+
delay = asyncRetryDelay
|
| 664 |
+
} else {
|
| 665 |
+
delay = asyncDebounce
|
| 666 |
+
}
|
| 667 |
}
|
| 668 |
+
}()
|
| 669 |
+
}
|
| 670 |
+
|
| 671 |
+
// flush 把待写内容落盘。必须在未持有 s.mu 的情况下调用。
|
| 672 |
+
func (s *Store) flush() error {
|
| 673 |
+
if s == nil {
|
| 674 |
+
return errors.New("chat history store is nil")
|
| 675 |
}
|
| 676 |
+
s.writeMu.Lock()
|
| 677 |
+
defer s.writeMu.Unlock()
|
| 678 |
+
|
| 679 |
+
s.mu.Lock()
|
| 680 |
+
if s.err != nil {
|
| 681 |
+
err := s.err
|
| 682 |
+
s.mu.Unlock()
|
| 683 |
+
return err
|
| 684 |
+
}
|
| 685 |
+
if !s.hasPendingLocked() {
|
| 686 |
+
s.mu.Unlock()
|
| 687 |
+
return nil
|
| 688 |
+
}
|
| 689 |
+
dirtyIDs := s.dirty
|
| 690 |
+
progressIDs := s.progress
|
| 691 |
+
deletedIDs := s.deleted
|
| 692 |
+
s.dirty = map[string]struct{}{}
|
| 693 |
+
s.progress = map[string]struct{}{}
|
| 694 |
+
s.deleted = map[string]struct{}{}
|
| 695 |
+
s.indexDirty = false
|
| 696 |
+
index := cloneFile(s.state)
|
| 697 |
+
fullEntries := make(map[string]Entry, len(dirtyIDs))
|
| 698 |
+
for id := range dirtyIDs {
|
| 699 |
+
if item, ok := s.details[id]; ok {
|
| 700 |
+
fullEntries[id] = cloneEntry(item)
|
| 701 |
+
}
|
| 702 |
+
}
|
| 703 |
+
progressEntries := make(map[string]Entry, len(progressIDs))
|
| 704 |
+
for id := range progressIDs {
|
| 705 |
+
if item, ok := s.details[id]; ok {
|
| 706 |
+
progressEntries[id] = cloneEntry(item)
|
| 707 |
+
}
|
| 708 |
+
}
|
| 709 |
+
s.mu.Unlock()
|
| 710 |
+
|
| 711 |
+
err := s.writeSnapshot(index, fullEntries, progressEntries, deletedIDs)
|
| 712 |
+
if err != nil {
|
| 713 |
+
s.restorePending(dirtyIDs, progressIDs, deletedIDs)
|
| 714 |
+
}
|
| 715 |
+
return err
|
| 716 |
+
}
|
| 717 |
+
|
| 718 |
+
// restorePending 在写盘失败后把待写标记放回去,让后续写入重试,
|
| 719 |
+
// 同时不覆盖期间新产生的标记。
|
| 720 |
+
func (s *Store) restorePending(dirtyIDs, progressIDs, deletedIDs map[string]struct{}) {
|
| 721 |
+
s.mu.Lock()
|
| 722 |
+
defer s.mu.Unlock()
|
| 723 |
+
for id := range dirtyIDs {
|
| 724 |
+
s.dirty[id] = struct{}{}
|
| 725 |
+
delete(s.progress, id)
|
| 726 |
+
delete(s.deleted, id)
|
| 727 |
+
}
|
| 728 |
+
for id := range progressIDs {
|
| 729 |
+
if _, isFull := s.dirty[id]; isFull {
|
| 730 |
continue
|
| 731 |
}
|
| 732 |
+
s.progress[id] = struct{}{}
|
| 733 |
+
delete(s.deleted, id)
|
| 734 |
+
}
|
| 735 |
+
for id := range deletedIDs {
|
| 736 |
+
s.deleted[id] = struct{}{}
|
| 737 |
+
delete(s.dirty, id)
|
| 738 |
+
delete(s.progress, id)
|
| 739 |
+
}
|
| 740 |
+
s.indexDirty = true
|
| 741 |
+
}
|
| 742 |
+
|
| 743 |
+
func (s *Store) detailPath(id string) string {
|
| 744 |
+
return filepath.Join(s.detailDir, strings.TrimSpace(id)+".json")
|
| 745 |
+
}
|
| 746 |
+
|
| 747 |
+
func (s *Store) livePath(id string) string {
|
| 748 |
+
return filepath.Join(s.detailDir, strings.TrimSpace(id)+".live.json")
|
| 749 |
+
}
|
| 750 |
+
|
| 751 |
+
func (s *Store) writeSnapshot(index File, fullEntries, progressEntries map[string]Entry, deletedIDs map[string]struct{}) error {
|
| 752 |
+
if err := os.MkdirAll(s.detailDir, 0o755); err != nil {
|
| 753 |
+
return fmt.Errorf("create chat history detail dir: %w", err)
|
| 754 |
+
}
|
| 755 |
+
for _, id := range sortedDetailIDs(deletedIDs) {
|
| 756 |
+
for _, path := range []string{s.detailPath(id), s.livePath(id)} {
|
| 757 |
+
if err := os.Remove(path); err != nil && !errors.Is(err, os.ErrNotExist) {
|
| 758 |
+
return fmt.Errorf("remove stale chat history detail: %w", err)
|
| 759 |
+
}
|
| 760 |
+
}
|
| 761 |
+
}
|
| 762 |
+
for _, id := range sortedDetailIDs(entryIDSet(fullEntries)) {
|
| 763 |
payload, err := json.MarshalIndent(detailEnvelope{
|
| 764 |
Version: FileVersion,
|
| 765 |
+
Item: fullEntries[id],
|
| 766 |
}, "", " ")
|
| 767 |
if err != nil {
|
| 768 |
return fmt.Errorf("encode chat history detail: %w", err)
|
| 769 |
}
|
| 770 |
+
if err := writeFileAtomic(s.detailPath(id), append(payload, '\n')); err != nil {
|
| 771 |
+
return err
|
| 772 |
+
}
|
| 773 |
+
// 完整文件已经带上最新状态,增量文件不再需要。
|
| 774 |
+
if err := os.Remove(s.livePath(id)); err != nil && !errors.Is(err, os.ErrNotExist) {
|
| 775 |
+
return fmt.Errorf("remove stale chat history progress: %w", err)
|
| 776 |
+
}
|
| 777 |
+
}
|
| 778 |
+
for _, id := range sortedDetailIDs(entryIDSet(progressEntries)) {
|
| 779 |
+
payload, err := json.MarshalIndent(progressEnvelope{
|
| 780 |
+
Version: FileVersion,
|
| 781 |
+
Item: progressFromEntry(progressEntries[id]),
|
| 782 |
+
}, "", " ")
|
| 783 |
+
if err != nil {
|
| 784 |
+
return fmt.Errorf("encode chat history progress: %w", err)
|
| 785 |
+
}
|
| 786 |
+
if err := writeFileAtomic(s.livePath(id), append(payload, '\n')); err != nil {
|
| 787 |
return err
|
| 788 |
}
|
| 789 |
}
|
| 790 |
|
| 791 |
+
payload, err := json.MarshalIndent(index, "", " ")
|
| 792 |
if err != nil {
|
| 793 |
return fmt.Errorf("encode chat history index: %w", err)
|
| 794 |
}
|
| 795 |
if err := writeFileAtomic(s.path, append(payload, '\n')); err != nil {
|
| 796 |
return err
|
| 797 |
}
|
|
|
|
| 798 |
return nil
|
| 799 |
}
|
| 800 |
|
| 801 |
+
func (s *Store) readEntryLocked(id string) (Entry, error) {
|
| 802 |
+
item, err := readDetailFile(s.detailPath(id))
|
| 803 |
+
if err != nil {
|
| 804 |
+
return Entry{}, err
|
| 805 |
+
}
|
| 806 |
+
if progress, ok := readProgressFile(s.livePath(id)); ok && progress.Revision > item.Revision {
|
| 807 |
+
progress.applyTo(&item)
|
| 808 |
+
}
|
| 809 |
+
return item, nil
|
| 810 |
+
}
|
| 811 |
+
|
| 812 |
func (s *Store) rebuildIndexLocked() {
|
| 813 |
summaries := make([]SummaryEntry, 0, len(s.details))
|
| 814 |
for _, item := range s.details {
|
|
|
|
| 852 |
next = s.state.Revision + 1
|
| 853 |
}
|
| 854 |
s.state.Revision = next
|
| 855 |
+
s.indexDirty = true
|
| 856 |
return next
|
| 857 |
}
|
| 858 |
|
|
|
|
| 895 |
return candidate
|
| 896 |
}
|
| 897 |
|
| 898 |
+
func progressFromEntry(item Entry) progressRecord {
|
| 899 |
+
return progressRecord{
|
| 900 |
+
Version: FileVersion,
|
| 901 |
+
Revision: item.Revision,
|
| 902 |
+
UpdatedAt: item.UpdatedAt,
|
| 903 |
+
CompletedAt: item.CompletedAt,
|
| 904 |
+
Status: item.Status,
|
| 905 |
+
ReasoningContent: item.ReasoningContent,
|
| 906 |
+
Content: item.Content,
|
| 907 |
+
Error: item.Error,
|
| 908 |
+
StatusCode: item.StatusCode,
|
| 909 |
+
ElapsedMs: item.ElapsedMs,
|
| 910 |
+
FinishReason: item.FinishReason,
|
| 911 |
+
Usage: cloneMap(item.Usage),
|
| 912 |
+
}
|
| 913 |
+
}
|
| 914 |
+
|
| 915 |
+
func (r progressRecord) applyTo(item *Entry) {
|
| 916 |
+
if item == nil {
|
| 917 |
+
return
|
| 918 |
+
}
|
| 919 |
+
item.Revision = r.Revision
|
| 920 |
+
item.UpdatedAt = r.UpdatedAt
|
| 921 |
+
item.CompletedAt = r.CompletedAt
|
| 922 |
+
item.Status = r.Status
|
| 923 |
+
item.ReasoningContent = r.ReasoningContent
|
| 924 |
+
item.Content = r.Content
|
| 925 |
+
item.Error = r.Error
|
| 926 |
+
item.StatusCode = r.StatusCode
|
| 927 |
+
item.ElapsedMs = r.ElapsedMs
|
| 928 |
+
item.FinishReason = r.FinishReason
|
| 929 |
+
item.Usage = cloneMap(r.Usage)
|
| 930 |
+
}
|
| 931 |
+
|
| 932 |
func readDetailFile(path string) (Entry, error) {
|
| 933 |
raw, err := os.ReadFile(path)
|
| 934 |
if err != nil {
|
|
|
|
| 941 |
return cloneEntry(env.Item), nil
|
| 942 |
}
|
| 943 |
|
| 944 |
+
// readProgressFile 读取流式增量文件。它只是完整文件之上的一个覆盖层,
|
| 945 |
+
// 缺失或损坏都不影响完整文件本身,因此这里不返回错误。
|
| 946 |
+
func readProgressFile(path string) (progressRecord, bool) {
|
| 947 |
+
raw, err := os.ReadFile(path)
|
| 948 |
+
if err != nil {
|
| 949 |
+
return progressRecord{}, false
|
| 950 |
+
}
|
| 951 |
+
var env progressEnvelope
|
| 952 |
+
if err := json.Unmarshal(raw, &env); err != nil {
|
| 953 |
+
config.Logger.Warn("[chat_history] ignoring corrupt progress file", "path", path, "error", err)
|
| 954 |
+
return progressRecord{}, false
|
| 955 |
+
}
|
| 956 |
+
return env.Item, true
|
| 957 |
+
}
|
| 958 |
+
|
| 959 |
func parseLegacy(raw []byte) (legacyFile, bool, error) {
|
| 960 |
var legacy legacyFile
|
| 961 |
if err := json.Unmarshal(raw, &legacy); err != nil {
|
|
|
|
| 975 |
return legacy, true, nil
|
| 976 |
}
|
| 977 |
|
| 978 |
+
// writeFileAtomic 用「临时文件 + rename」保证读到的永远是一个完整文件。
|
| 979 |
+
//
|
| 980 |
+
// 这里刻意不做 fsync:chat history 是可再生的诊断数据,而本项目的主要部署
|
| 981 |
+
// 目标(HuggingFace Spaces 的 bucket 卷等对象存储挂载)单次 fsync 的延迟远高于
|
| 982 |
+
// 本地盘。流式响应期间每个分片都 fsync 会把整个进程拖垮,收益却只是「掉电时
|
| 983 |
+
// 少丢一条历史」。rename 本身已经保证不会出现半截 JSON。
|
| 984 |
func writeFileAtomic(path string, body []byte) error {
|
| 985 |
dir := filepath.Dir(path)
|
| 986 |
if dir == "" {
|
|
|
|
| 1015 |
if _, err := tmpFile.Write(body); err != nil {
|
| 1016 |
return withCleanup(fmt.Errorf("write temp chat history: %w", err), tmpFile.Close())
|
| 1017 |
}
|
|
|
|
|
|
|
|
|
|
| 1018 |
if err := tmpFile.Close(); err != nil {
|
| 1019 |
if cleanupErr := cleanup(); cleanupErr != nil {
|
| 1020 |
return errors.Join(fmt.Errorf("close temp chat history: %w", err), cleanupErr)
|
|
|
|
| 1051 |
if s.dirty == nil {
|
| 1052 |
s.dirty = map[string]struct{}{}
|
| 1053 |
}
|
| 1054 |
+
if s.progress == nil {
|
| 1055 |
+
s.progress = map[string]struct{}{}
|
| 1056 |
+
}
|
| 1057 |
if s.deleted == nil {
|
| 1058 |
s.deleted = map[string]struct{}{}
|
| 1059 |
}
|
| 1060 |
s.dirty[id] = struct{}{}
|
| 1061 |
+
delete(s.progress, id)
|
| 1062 |
+
delete(s.deleted, id)
|
| 1063 |
+
}
|
| 1064 |
+
|
| 1065 |
+
// markProgressDirtyLocked 只安排一次增量写入。若该条目已经排了完整写入,
|
| 1066 |
+
// 完整写入本身就包含最新状态,不再额外写增量文件。
|
| 1067 |
+
func (s *Store) markProgressDirtyLocked(id string) {
|
| 1068 |
+
id = strings.TrimSpace(id)
|
| 1069 |
+
if id == "" {
|
| 1070 |
+
return
|
| 1071 |
+
}
|
| 1072 |
+
if s.progress == nil {
|
| 1073 |
+
s.progress = map[string]struct{}{}
|
| 1074 |
+
}
|
| 1075 |
+
if s.deleted == nil {
|
| 1076 |
+
s.deleted = map[string]struct{}{}
|
| 1077 |
+
}
|
| 1078 |
+
if _, isFull := s.dirty[id]; isFull {
|
| 1079 |
+
return
|
| 1080 |
+
}
|
| 1081 |
+
s.progress[id] = struct{}{}
|
| 1082 |
delete(s.deleted, id)
|
| 1083 |
}
|
| 1084 |
|
|
|
|
| 1090 |
if s.dirty == nil {
|
| 1091 |
s.dirty = map[string]struct{}{}
|
| 1092 |
}
|
| 1093 |
+
if s.progress == nil {
|
| 1094 |
+
s.progress = map[string]struct{}{}
|
| 1095 |
+
}
|
| 1096 |
if s.deleted == nil {
|
| 1097 |
s.deleted = map[string]struct{}{}
|
| 1098 |
}
|
| 1099 |
s.deleted[id] = struct{}{}
|
| 1100 |
delete(s.dirty, id)
|
| 1101 |
+
delete(s.progress, id)
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1102 |
}
|
| 1103 |
|
| 1104 |
func sortedDetailIDs(ids map[string]struct{}) []string {
|
|
|
|
| 1113 |
return out
|
| 1114 |
}
|
| 1115 |
|
| 1116 |
+
func entryIDSet(entries map[string]Entry) map[string]struct{} {
|
| 1117 |
+
if len(entries) == 0 {
|
| 1118 |
+
return nil
|
| 1119 |
+
}
|
| 1120 |
+
out := make(map[string]struct{}, len(entries))
|
| 1121 |
+
for id := range entries {
|
| 1122 |
+
out[id] = struct{}{}
|
| 1123 |
+
}
|
| 1124 |
+
return out
|
| 1125 |
+
}
|
| 1126 |
+
|
| 1127 |
func cloneFile(in File) File {
|
| 1128 |
out := File{
|
| 1129 |
Version: in.Version,
|
internal/chathistory/store_test.go
CHANGED
|
@@ -3,6 +3,7 @@ package chathistory
|
|
| 3 |
import (
|
| 4 |
"bytes"
|
| 5 |
"encoding/json"
|
|
|
|
| 6 |
"os"
|
| 7 |
"path/filepath"
|
| 8 |
"strings"
|
|
@@ -633,3 +634,163 @@ func TestUpdateAllowsOverwritingContentWithNewValue(t *testing.T) {
|
|
| 633 |
t.Fatalf("expected content to be overwritten, got %q", updated.Content)
|
| 634 |
}
|
| 635 |
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 3 |
import (
|
| 4 |
"bytes"
|
| 5 |
"encoding/json"
|
| 6 |
+
"errors"
|
| 7 |
"os"
|
| 8 |
"path/filepath"
|
| 9 |
"strings"
|
|
|
|
| 634 |
t.Fatalf("expected content to be overwritten, got %q", updated.Content)
|
| 635 |
}
|
| 636 |
}
|
| 637 |
+
|
| 638 |
+
func waitForFile(t *testing.T, path string, timeout time.Duration) []byte {
|
| 639 |
+
t.Helper()
|
| 640 |
+
deadline := time.Now().Add(timeout)
|
| 641 |
+
for time.Now().Before(deadline) {
|
| 642 |
+
if raw, err := os.ReadFile(path); err == nil {
|
| 643 |
+
return raw
|
| 644 |
+
}
|
| 645 |
+
time.Sleep(5 * time.Millisecond)
|
| 646 |
+
}
|
| 647 |
+
t.Fatalf("timed out waiting for %s", path)
|
| 648 |
+
return nil
|
| 649 |
+
}
|
| 650 |
+
|
| 651 |
+
// 流式进度必须写「小增量文件」,而不是每次重写包含完整上下文的大文件:
|
| 652 |
+
// 这是 HuggingFace Spaces 等对象存储挂载上拖慢整个服务的主因。
|
| 653 |
+
func TestStoreBackgroundProgressUsesCompactOverlay(t *testing.T) {
|
| 654 |
+
path := filepath.Join(t.TempDir(), "chat_history.json")
|
| 655 |
+
store := New(path)
|
| 656 |
+
t.Cleanup(store.Close)
|
| 657 |
+
|
| 658 |
+
// 结尾不能是空白:Start 会对 FinalPrompt 做 TrimSpace。
|
| 659 |
+
bigContext := strings.Repeat("context line\n", 2000) + "tail"
|
| 660 |
+
started, err := store.Start(StartParams{
|
| 661 |
+
CallerID: "caller:abc",
|
| 662 |
+
Model: "deepseek-v4-flash",
|
| 663 |
+
Stream: true,
|
| 664 |
+
UserInput: "hello",
|
| 665 |
+
Messages: []Message{{Role: "user", Content: bigContext}},
|
| 666 |
+
HistoryText: bigContext,
|
| 667 |
+
FinalPrompt: bigContext,
|
| 668 |
+
})
|
| 669 |
+
if err != nil {
|
| 670 |
+
t.Fatalf("start failed: %v", err)
|
| 671 |
+
}
|
| 672 |
+
|
| 673 |
+
fullPath := filepath.Join(store.DetailDir(), started.ID+".json")
|
| 674 |
+
livePath := filepath.Join(store.DetailDir(), started.ID+".live.json")
|
| 675 |
+
beforeFull, err := os.ReadFile(fullPath)
|
| 676 |
+
if err != nil {
|
| 677 |
+
t.Fatalf("read full detail failed: %v", err)
|
| 678 |
+
}
|
| 679 |
+
|
| 680 |
+
if _, err := store.Update(started.ID, UpdateParams{
|
| 681 |
+
Status: "streaming",
|
| 682 |
+
ReasoningContent: "let me think",
|
| 683 |
+
Content: "partial answer",
|
| 684 |
+
ElapsedMs: 12,
|
| 685 |
+
Background: true,
|
| 686 |
+
}); err != nil {
|
| 687 |
+
t.Fatalf("background update failed: %v", err)
|
| 688 |
+
}
|
| 689 |
+
|
| 690 |
+
liveRaw := waitForFile(t, livePath, 5*time.Second)
|
| 691 |
+
|
| 692 |
+
afterFull, err := os.ReadFile(fullPath)
|
| 693 |
+
if err != nil {
|
| 694 |
+
t.Fatalf("re-read full detail failed: %v", err)
|
| 695 |
+
}
|
| 696 |
+
if !bytes.Equal(beforeFull, afterFull) {
|
| 697 |
+
t.Fatalf("background progress must not rewrite the full detail file")
|
| 698 |
+
}
|
| 699 |
+
if len(liveRaw) >= len(beforeFull) {
|
| 700 |
+
t.Fatalf("expected overlay (%d bytes) to be far smaller than full detail (%d bytes)", len(liveRaw), len(beforeFull))
|
| 701 |
+
}
|
| 702 |
+
t.Logf("progress write: overlay=%d bytes vs full detail=%d bytes (%.1fx smaller)",
|
| 703 |
+
len(liveRaw), len(beforeFull), float64(len(beforeFull))/float64(len(liveRaw)))
|
| 704 |
+
|
| 705 |
+
reloaded := New(path)
|
| 706 |
+
t.Cleanup(reloaded.Close)
|
| 707 |
+
full, err := reloaded.Get(started.ID)
|
| 708 |
+
if err != nil {
|
| 709 |
+
t.Fatalf("get after reload failed: %v", err)
|
| 710 |
+
}
|
| 711 |
+
if full.Content != "partial answer" || full.ReasoningContent != "let me think" || full.Status != "streaming" {
|
| 712 |
+
t.Fatalf("expected overlay to be applied on load, got %#v", full)
|
| 713 |
+
}
|
| 714 |
+
if full.HistoryText != bigContext || full.FinalPrompt != bigContext {
|
| 715 |
+
t.Fatalf("expected immutable context to survive reload")
|
| 716 |
+
}
|
| 717 |
+
if len(full.Messages) != 1 || full.Messages[0].Content != bigContext {
|
| 718 |
+
t.Fatalf("expected messages to survive reload")
|
| 719 |
+
}
|
| 720 |
+
}
|
| 721 |
+
|
| 722 |
+
// 终态更新必须回写成完整文件并清理增量文件,避免残留覆盖层。
|
| 723 |
+
func TestStoreTerminalUpdateReplacesProgressOverlay(t *testing.T) {
|
| 724 |
+
path := filepath.Join(t.TempDir(), "chat_history.json")
|
| 725 |
+
store := New(path)
|
| 726 |
+
t.Cleanup(store.Close)
|
| 727 |
+
|
| 728 |
+
started, err := store.Start(StartParams{UserInput: "hello", HistoryText: "ctx", FinalPrompt: "ctx"})
|
| 729 |
+
if err != nil {
|
| 730 |
+
t.Fatalf("start failed: %v", err)
|
| 731 |
+
}
|
| 732 |
+
livePath := filepath.Join(store.DetailDir(), started.ID+".live.json")
|
| 733 |
+
|
| 734 |
+
if _, err := store.Update(started.ID, UpdateParams{
|
| 735 |
+
Status: "streaming",
|
| 736 |
+
Content: "partial",
|
| 737 |
+
Background: true,
|
| 738 |
+
}); err != nil {
|
| 739 |
+
t.Fatalf("background update failed: %v", err)
|
| 740 |
+
}
|
| 741 |
+
waitForFile(t, livePath, 5*time.Second)
|
| 742 |
+
|
| 743 |
+
if _, err := store.Update(started.ID, UpdateParams{
|
| 744 |
+
Status: "success",
|
| 745 |
+
Content: "final answer",
|
| 746 |
+
Completed: true,
|
| 747 |
+
}); err != nil {
|
| 748 |
+
t.Fatalf("terminal update failed: %v", err)
|
| 749 |
+
}
|
| 750 |
+
if _, err := os.Stat(livePath); !errors.Is(err, os.ErrNotExist) {
|
| 751 |
+
t.Fatalf("expected overlay to be removed after terminal update, stat err=%v", err)
|
| 752 |
+
}
|
| 753 |
+
|
| 754 |
+
reloaded := New(path)
|
| 755 |
+
t.Cleanup(reloaded.Close)
|
| 756 |
+
full, err := reloaded.Get(started.ID)
|
| 757 |
+
if err != nil {
|
| 758 |
+
t.Fatalf("get after reload failed: %v", err)
|
| 759 |
+
}
|
| 760 |
+
if full.Content != "final answer" || full.Status != "success" || full.CompletedAt == 0 {
|
| 761 |
+
t.Fatalf("expected terminal state to be persisted, got %#v", full)
|
| 762 |
+
}
|
| 763 |
+
}
|
| 764 |
+
|
| 765 |
+
// 流式进度的落盘失败不能冒泡给调用方,否则会打断正在进行的 SSE 响应;
|
| 766 |
+
// 同时内存状态必须保持可读。
|
| 767 |
+
func TestStoreBackgroundProgressDoesNotSurfaceWriteFailure(t *testing.T) {
|
| 768 |
+
path := filepath.Join(t.TempDir(), "chat_history.json")
|
| 769 |
+
store := New(path)
|
| 770 |
+
t.Cleanup(store.Close)
|
| 771 |
+
|
| 772 |
+
started, err := store.Start(StartParams{UserInput: "hello"})
|
| 773 |
+
if err != nil {
|
| 774 |
+
t.Fatalf("start failed: %v", err)
|
| 775 |
+
}
|
| 776 |
+
restore := blockDetailDir(t, store.DetailDir())
|
| 777 |
+
defer restore()
|
| 778 |
+
|
| 779 |
+
if _, err := store.Update(started.ID, UpdateParams{
|
| 780 |
+
Status: "streaming",
|
| 781 |
+
Content: "partial",
|
| 782 |
+
Background: true,
|
| 783 |
+
}); err != nil {
|
| 784 |
+
t.Fatalf("background progress must not surface write errors, got %v", err)
|
| 785 |
+
}
|
| 786 |
+
if err := store.Err(); err != nil {
|
| 787 |
+
t.Fatalf("transient background write failure should not latch store error: %v", err)
|
| 788 |
+
}
|
| 789 |
+
snapshot, err := store.Snapshot()
|
| 790 |
+
if err != nil {
|
| 791 |
+
t.Fatalf("snapshot must keep working while writes fail: %v", err)
|
| 792 |
+
}
|
| 793 |
+
if len(snapshot.Items) != 1 {
|
| 794 |
+
t.Fatalf("expected in-memory state to stay readable, got %#v", snapshot.Items)
|
| 795 |
+
}
|
| 796 |
+
}
|
internal/httpapi/openai/chat/chat_history.go
CHANGED
|
@@ -132,12 +132,15 @@ func (s *chatHistorySession) progress(thinking, content string) {
|
|
| 132 |
return
|
| 133 |
}
|
| 134 |
s.lastPersist = now
|
|
|
|
|
|
|
| 135 |
s.persistUpdate(chathistory.UpdateParams{
|
| 136 |
Status: "streaming",
|
| 137 |
ReasoningContent: thinking,
|
| 138 |
Content: content,
|
| 139 |
StatusCode: http.StatusOK,
|
| 140 |
ElapsedMs: time.Since(s.startedAt).Milliseconds(),
|
|
|
|
| 141 |
})
|
| 142 |
}
|
| 143 |
|
|
|
|
| 132 |
return
|
| 133 |
}
|
| 134 |
s.lastPersist = now
|
| 135 |
+
// Background=true:流式进度只写增量文件并由后台合并落盘,避免在 SSE
|
| 136 |
+
// 读取循环里同步等待磁盘(对象存储挂载下单次写入延迟很高)。
|
| 137 |
s.persistUpdate(chathistory.UpdateParams{
|
| 138 |
Status: "streaming",
|
| 139 |
ReasoningContent: thinking,
|
| 140 |
Content: content,
|
| 141 |
StatusCode: http.StatusOK,
|
| 142 |
ElapsedMs: time.Since(s.startedAt).Milliseconds(),
|
| 143 |
+
Background: true,
|
| 144 |
})
|
| 145 |
}
|
| 146 |
|
internal/responsehistory/session.go
CHANGED
|
@@ -127,12 +127,15 @@ func (s *Session) Progress(thinking, content string) {
|
|
| 127 |
return
|
| 128 |
}
|
| 129 |
s.lastPersist = now
|
|
|
|
|
|
|
| 130 |
s.persistUpdate(chathistory.UpdateParams{
|
| 131 |
Status: "streaming",
|
| 132 |
ReasoningContent: thinking,
|
| 133 |
Content: content,
|
| 134 |
StatusCode: http.StatusOK,
|
| 135 |
ElapsedMs: time.Since(s.startedAt).Milliseconds(),
|
|
|
|
| 136 |
})
|
| 137 |
}
|
| 138 |
|
|
|
|
| 127 |
return
|
| 128 |
}
|
| 129 |
s.lastPersist = now
|
| 130 |
+
// Background=true:流式进度只写增量文件并由后台合并落盘,避免在 SSE
|
| 131 |
+
// 读取循环里同步等待磁盘(对象存储挂载下单次写入延迟很高)。
|
| 132 |
s.persistUpdate(chathistory.UpdateParams{
|
| 133 |
Status: "streaming",
|
| 134 |
ReasoningContent: thinking,
|
| 135 |
Content: content,
|
| 136 |
StatusCode: http.StatusOK,
|
| 137 |
ElapsedMs: time.Since(s.startedAt).Milliseconds(),
|
| 138 |
+
Background: true,
|
| 139 |
})
|
| 140 |
}
|
| 141 |
|
webui/src/features/chatHistory/ChatHistoryContainer.jsx
CHANGED
|
@@ -10,8 +10,11 @@ import {
|
|
| 10 |
VIEW_MODE_KEY,
|
| 11 |
} from './chatHistoryUtils'
|
| 12 |
|
| 13 |
-
|
| 14 |
-
|
|
|
|
|
|
|
|
|
|
| 15 |
|
| 16 |
export default function ChatHistoryContainer({ authFetch, onMessage }) {
|
| 17 |
const { t, lang } = useI18n()
|
|
|
|
| 10 |
VIEW_MODE_KEY,
|
| 11 |
} from './chatHistoryUtils'
|
| 12 |
|
| 13 |
+
// 列表 / 详情轮询间隔。流式期间全局 revision 每 250ms 就变一次,ETag 永远
|
| 14 |
+
// 命中不了 304,所以每次轮询都是全量传输;在 HF Spaces 这类低配实例上,
|
| 15 |
+
// 过密的轮询本身会明显加重服务端负担,这里取一个仍然「够实时」的间隔。
|
| 16 |
+
const LIST_REFRESH_MS = 3000
|
| 17 |
+
const STREAMING_DETAIL_REFRESH_MS = 1500
|
| 18 |
|
| 19 |
export default function ChatHistoryContainer({ authFetch, onMessage }) {
|
| 20 |
const { t, lang } = useI18n()
|