Files
kenn-io--agentsview/internal/sync/rebuild_contributor_test.go
wehub-resource-sync f99010fae1
CI / lint (push) Failing after 1s
CI / frontend (push) Failing after 1s
CI / scripts (push) Failing after 1s
CI / Go Test (ubuntu-latest) (push) Failing after 0s
CI / frontend-node-25 (push) Failing after 1s
CI / docs (push) Failing after 0s
CI / coverage (push) Failing after 0s
CI / e2e (push) Failing after 0s
Docker / build-and-push (push) Failing after 1s
CI / integration (push) Failing after 4m43s
CI / Go Test (windows-latest) (push) Has been cancelled
CI / Desktop Unit Tests (Windows) (push) Has been cancelled
Desktop Artifacts / Desktop Build (Linux (arm64)) (push) Has been cancelled
Desktop Artifacts / Desktop Build (Linux) (push) Has been cancelled
Desktop Artifacts / Desktop Build (Windows) (push) Has been cancelled
Desktop Artifacts (macOS) / Desktop Build (macOS (aarch64)) (push) Has been cancelled
Desktop Artifacts (macOS) / Desktop Build (macOS (x86_64)) (push) Has been cancelled
chore: import upstream snapshot with attribution
2026-07-13 12:30:36 +08:00

809 lines
28 KiB
Go

package sync
import (
"context"
"encoding/json"
"errors"
"fmt"
"os"
"path/filepath"
"testing"
"time"
"go.kenn.io/agentsview/internal/db"
"go.kenn.io/agentsview/internal/parser"
"go.kenn.io/agentsview/internal/testjsonl"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestMergeSyncStatsIncludesAdditiveRebuildFields(t *testing.T) {
dst := SyncStats{
OrphanedCopied: 2,
RebuildPhases: []RebuildPhaseStats{{Contributor: "local", BatchedWrites: 1}},
}
src := SyncStats{
OrphanedCopied: 3,
RebuildPhases: []RebuildPhaseStats{{Contributor: "remote", BatchedWrites: 2}},
}
mergeSyncStats(&dst, src)
assert.Equal(t, 5, dst.OrphanedCopied)
assert.Equal(t, []RebuildPhaseStats{
{Contributor: "local", BatchedWrites: 1},
{Contributor: "remote", BatchedWrites: 2},
}, dst.RebuildPhases)
}
func TestResyncAllLegacyOmitsContributorDiagnostics(t *testing.T) {
root := t.TempDir()
database, err := db.Open(filepath.Join(t.TempDir(), "archive.db"))
require.NoError(t, err)
t.Cleanup(func() { require.NoError(t, database.Close()) })
engine := NewEngine(database, EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentClaude: {root}},
Machine: "local",
})
t.Cleanup(engine.Close)
path := filepath.Join(root, "project", "legacy.jsonl")
require.NoError(t, os.MkdirAll(filepath.Dir(path), 0o755))
require.NoError(t, os.WriteFile(path, []byte(testjsonl.NewSessionBuilder().
AddClaudeUser("2026-01-01T00:00:00Z", "legacy phase fixture").String()), 0o644))
stats := engine.ResyncAll(context.Background(), nil)
require.False(t, stats.Aborted)
assert.Nil(t, stats.RebuildPhases)
assert.Nil(t, engine.LastSyncStats().RebuildPhases)
payload, err := json.Marshal(stats)
require.NoError(t, err)
assert.NotContains(t, string(payload), "rebuild_phases")
}
func TestResyncContributorsRunInOrderWithCumulativeProgress(t *testing.T) {
localRoot := t.TempDir()
rootA := t.TempDir()
rootB := t.TempDir()
for _, fixture := range []struct {
root string
project string
id string
content string
}{
{localRoot, "local", "local", "local cumulative progress"},
{rootA, "a", "contributor-a", "contributor A progress"},
{rootB, "b", "contributor-b", "contributor B progress"},
} {
path := filepath.Join(fixture.root, fixture.project, fixture.id+".jsonl")
require.NoError(t, os.MkdirAll(filepath.Dir(path), 0o755))
require.NoError(t, os.WriteFile(path, []byte(testjsonl.NewSessionBuilder().
AddClaudeUser("2026-01-01T00:00:00Z", fixture.content).String()), 0o644))
}
database, err := db.Open(filepath.Join(t.TempDir(), "archive.db"))
require.NoError(t, err)
t.Cleanup(func() { require.NoError(t, database.Close()) })
engine := NewEngine(database, EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentClaude: {localRoot}},
Machine: "local",
})
t.Cleanup(engine.Close)
order := []string{}
progressByContributor := map[string]Progress{}
labelProgress := func(name string) func(Progress) Progress {
return func(p Progress) Progress {
p.Detail = name
return p
}
}
onProgress := func(p Progress) {
if p.Detail == "A" || p.Detail == "B" {
progressByContributor[p.Detail] = p
}
}
ftsCalls := 0
stats, err := engine.resyncAllWithOptionsAndOperations(
context.Background(), onProgress, RebuildOptions{
Contributors: []RebuildContributor{
{
Name: "A",
Config: EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentClaude: {rootA}},
Machine: "A", IDPrefix: "A~", Ephemeral: true,
},
Progress: labelProgress("A"),
AfterSync: func(*Engine, *db.DB) error {
order = append(order, "A")
return nil
},
},
{
Name: "B",
Config: EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentClaude: {rootB}},
Machine: "B", IDPrefix: "B~", Ephemeral: true,
},
Progress: labelProgress("B"),
AfterSync: func(*Engine, *db.DB) error {
order = append(order, "B")
return nil
},
},
},
}, rebuildOperations{
rebuildFTS: func(database *db.DB) error {
ftsCalls++
return database.RebuildFTS()
},
},
)
require.NoError(t, err)
require.False(t, stats.Aborted, "resync aborted: %+v", stats)
assert.Equal(t, []string{"A", "B"}, order)
assert.Equal(t, 3, stats.Synced)
assert.Equal(t, []string{"local", "A", "B"}, []string{
stats.RebuildPhases[0].Contributor,
stats.RebuildPhases[1].Contributor,
stats.RebuildPhases[2].Contributor,
})
assert.Equal(t, 1, ftsCalls, "FTS rebuilt once instead of per contributor")
assert.Equal(t, Progress{
Phase: PhaseDone, Detail: "A", Resync: true,
SessionsTotal: 2, SessionsDone: 2, MessagesIndexed: 2,
}, progressByContributor["A"])
assert.Equal(t, Progress{
Phase: PhaseDone, Detail: "B", Resync: true,
SessionsTotal: 3, SessionsDone: 3, MessagesIndexed: 3,
}, progressByContributor["B"])
}
func TestResyncAbortsWhenContributorLosesHistoricalSource(t *testing.T) {
localRoot := t.TempDir()
remoteRoot := t.TempDir()
writeSession := func(root, project, name, content string) string {
path := filepath.Join(root, project, name+".jsonl")
require.NoError(t, os.MkdirAll(filepath.Dir(path), 0o755))
require.NoError(t, os.WriteFile(path, []byte(testjsonl.NewSessionBuilder().
AddClaudeUser("2026-01-01T00:00:00Z", content).String()), 0o644))
return path
}
writeSession(localRoot, "local", "local", "local source")
remotePath := writeSession(remoteRoot, "remote", "remote", "remote source")
database, err := db.Open(filepath.Join(t.TempDir(), "archive.db"))
require.NoError(t, err)
t.Cleanup(func() { require.NoError(t, database.Close()) })
engine := NewEngine(database, EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentClaude: {localRoot}},
Machine: "local",
})
t.Cleanup(engine.Close)
options := RebuildOptions{Contributors: []RebuildContributor{{
Name: "remote",
Config: EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentClaude: {remoteRoot}},
Machine: "remote", IDPrefix: "remote~", Ephemeral: true,
},
}}}
initial, err := engine.ResyncAllWithOptions(context.Background(), nil, options)
require.NoError(t, err)
require.False(t, initial.Aborted)
require.NoError(t, os.Remove(remotePath))
stats, err := engine.ResyncAllWithOptions(context.Background(), nil, options)
require.NoError(t, err)
assert.True(t, stats.Aborted,
"a healthy local pass must not mask an empty historical contributor")
remote, err := database.GetSession(context.Background(), "remote~remote")
require.NoError(t, err)
assert.NotNil(t, remote, "aborted rebuild must preserve the active archive")
}
func TestResyncDoesNotTreatSameNamedContributorAsLocalHistory(t *testing.T) {
localRoot := t.TempDir()
remoteRoot := t.TempDir()
remotePath := filepath.Join(remoteRoot, "project", "remote.jsonl")
require.NoError(t, os.MkdirAll(filepath.Dir(remotePath), 0o755))
require.NoError(t, os.WriteFile(remotePath, []byte(testjsonl.NewSessionBuilder().
AddClaudeUser("2026-01-01T00:00:00Z", "same-name remote source").String()), 0o644))
database, err := db.Open(filepath.Join(t.TempDir(), "archive.db"))
require.NoError(t, err)
t.Cleanup(func() { require.NoError(t, database.Close()) })
engine := NewEngine(database, EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentClaude: {localRoot}},
Machine: "collector-host",
})
t.Cleanup(engine.Close)
options := RebuildOptions{Contributors: []RebuildContributor{{
Name: "collector-host",
Config: EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentClaude: {remoteRoot}},
Machine: "collector-host", IDPrefix: "collector-host~", Ephemeral: true,
},
}}}
initial, err := engine.ResyncAllWithOptions(context.Background(), nil, options)
require.NoError(t, err)
require.False(t, initial.Aborted)
stats, err := engine.ResyncAllWithOptions(context.Background(), nil, options)
require.NoError(t, err)
assert.False(t, stats.Aborted,
"remote-prefixed history must not make an empty local source look incomplete")
remote, err := database.GetSession(context.Background(), "collector-host~remote")
require.NoError(t, err)
assert.NotNil(t, remote)
}
func TestResyncContributorPostSwapReopenFailureReturnsCoordinatorError(t *testing.T) {
root := t.TempDir()
database, err := db.Open(filepath.Join(t.TempDir(), "archive.db"))
require.NoError(t, err)
t.Cleanup(func() { require.NoError(t, database.Close()) })
emitter := &fakeEmitter{}
engine := NewEngine(database, EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentClaude: {root}},
Machine: "local",
Emitter: emitter,
})
t.Cleanup(engine.Close)
path := filepath.Join(root, "project", "session.jsonl")
require.NoError(t, os.MkdirAll(filepath.Dir(path), 0o755))
require.NoError(t, os.WriteFile(path, []byte(testjsonl.NewSessionBuilder().
AddClaudeUser("2026-01-01T00:00:00Z", "post swap reopen fixture").String()), 0o644))
sentinel := errors.New("reopen sentinel")
stats, err := engine.resyncAllWithOptionsAndOperations(
context.Background(), nil, RebuildOptions{}, rebuildOperations{
rebuildFTS: productionRebuildOperations.rebuildFTS,
reopen: func(*db.DB) error { return sentinel },
},
)
require.ErrorIs(t, err, sentinel)
assert.True(t, stats.Aborted)
assert.Contains(t, stats.Warnings,
"resync swap completed but reopening active database failed: reopen sentinel")
assert.Empty(t, emitter.got(), "failed reopen published a successful sync")
assert.True(t, engine.LastSyncStats().Aborted)
// The rename already completed. Restore the handle and prove the rebuilt
// file is now the active archive even though publication failed.
require.NoError(t, database.Reopen())
session, getErr := database.GetSession(context.Background(), "session")
require.NoError(t, getErr)
require.NotNil(t, session)
}
func TestResyncContributorFTSFailureAbortsAndCleansTempDB(t *testing.T) {
root := t.TempDir()
database, err := db.Open(filepath.Join(t.TempDir(), "archive.db"))
require.NoError(t, err)
t.Cleanup(func() { require.NoError(t, database.Close()) })
engine := NewEngine(database, EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentClaude: {root}},
Machine: "local",
})
t.Cleanup(engine.Close)
oldPath := filepath.Join(root, "old", "old.jsonl")
require.NoError(t, os.MkdirAll(filepath.Dir(oldPath), 0o755))
require.NoError(t, os.WriteFile(oldPath, []byte(testjsonl.NewSessionBuilder().
AddClaudeUser("2026-01-01T00:00:00Z", "search survives failure").String()), 0o644))
require.Equal(t, 1, engine.SyncAll(context.Background(), nil).Synced)
require.NoError(t, os.Remove(oldPath))
newPath := filepath.Join(root, "new", "new.jsonl")
require.NoError(t, os.MkdirAll(filepath.Dir(newPath), 0o755))
require.NoError(t, os.WriteFile(newPath, []byte(testjsonl.NewSessionBuilder().
AddClaudeUser("2026-01-01T00:00:01Z", "partial new data").String()), 0o644))
sentinel := errors.New("fts sentinel")
engine.syncMu.Lock()
stats, err := engine.resyncAllWithOptionsLocked(
context.Background(), nil, RebuildOptions{}, rebuildOperations{
rebuildFTS: func(*db.DB) error { return sentinel },
},
)
engine.syncMu.Unlock()
require.ErrorIs(t, err, sentinel)
assert.True(t, stats.Aborted)
page, searchErr := database.Search(context.Background(), db.SearchFilter{
Query: "search survives failure", Limit: 5,
})
require.NoError(t, searchErr)
require.Len(t, page.Results, 1)
assert.NoFileExists(t, database.Path()+resyncTempSuffix)
assert.NoFileExists(t, database.Path()+resyncTempSuffix+"-wal")
assert.NoFileExists(t, database.Path()+resyncTempSuffix+"-shm")
}
type blockingRebuildProvider struct {
parser.ProviderBase
started chan struct{}
}
func (p *blockingRebuildProvider) Discover(context.Context) ([]parser.SourceRef, error) {
return []parser.SourceRef{{
Provider: parser.AgentCowork, Key: "blocked-source",
DisplayPath: "blocked-source", FingerprintKey: "blocked-source",
}}, nil
}
func (p *blockingRebuildProvider) Fingerprint(
context.Context, parser.SourceRef,
) (parser.SourceFingerprint, error) {
return parser.SourceFingerprint{Key: "blocked-source", Size: 1, MTimeNS: 1}, nil
}
func (p *blockingRebuildProvider) Parse(
ctx context.Context, _ parser.ParseRequest,
) (parser.ParseOutcome, error) {
close(p.started)
<-ctx.Done()
return parser.ParseOutcome{}, ctx.Err()
}
type blockingRebuildFactory struct{ provider *blockingRebuildProvider }
func (f blockingRebuildFactory) Definition() parser.AgentDef {
return f.provider.Definition()
}
func (f blockingRebuildFactory) Capabilities() parser.Capabilities {
return f.provider.Capabilities()
}
func (f blockingRebuildFactory) NewProvider(parser.ProviderConfig) parser.Provider {
return f.provider
}
type trackingRebuildProvider struct {
parser.ProviderBase
discoverCalls int
}
func (p *trackingRebuildProvider) Discover(context.Context) ([]parser.SourceRef, error) {
p.discoverCalls++
return nil, nil
}
func (p *trackingRebuildProvider) Parse(
context.Context, parser.ParseRequest,
) (parser.ParseOutcome, error) {
return parser.ParseOutcome{}, nil
}
type trackingRebuildFactory struct{ provider *trackingRebuildProvider }
func (f trackingRebuildFactory) Definition() parser.AgentDef {
return f.provider.Definition()
}
func (f trackingRebuildFactory) Capabilities() parser.Capabilities {
return f.provider.Capabilities()
}
func (f trackingRebuildFactory) NewProvider(parser.ProviderConfig) parser.Provider {
return f.provider
}
func TestResyncContributorCancellationPreservesArchiveAndCleansTempDB(t *testing.T) {
root := t.TempDir()
database, err := db.Open(filepath.Join(t.TempDir(), "archive.db"))
require.NoError(t, err)
t.Cleanup(func() { require.NoError(t, database.Close()) })
engine := NewEngine(database, EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentClaude: {root}},
Machine: "local",
})
t.Cleanup(engine.Close)
oldPath := filepath.Join(root, "old", "old.jsonl")
require.NoError(t, os.MkdirAll(filepath.Dir(oldPath), 0o755))
require.NoError(t, os.WriteFile(oldPath, []byte(testjsonl.NewSessionBuilder().
AddClaudeUser("2026-01-01T00:00:00Z", "archive before cancel").String()), 0o644))
require.Equal(t, 1, engine.SyncAll(context.Background(), nil).Synced)
require.NoError(t, os.Remove(oldPath))
provider := &blockingRebuildProvider{
ProviderBase: parser.ProviderBase{
Def: parser.AgentDef{Type: parser.AgentCowork},
Caps: parser.Capabilities{Source: parser.SourceCapabilities{
DiscoverSources: parser.CapabilitySupported,
}},
},
started: make(chan struct{}),
}
ctx, cancel := context.WithCancel(context.Background())
firstHookCalls := 0
secondHookCalls := 0
secondProvider := &trackingRebuildProvider{ProviderBase: parser.ProviderBase{
Def: parser.AgentDef{Type: parser.AgentCowork},
Caps: parser.Capabilities{Source: parser.SourceCapabilities{
DiscoverSources: parser.CapabilitySupported,
}},
}}
result := make(chan struct {
stats SyncStats
err error
}, 1)
go func() {
stats, runErr := engine.ResyncAllWithOptions(ctx, nil, RebuildOptions{
Contributors: []RebuildContributor{
{
Name: "blocking",
Config: EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentCowork: {root}},
Machine: "blocking", Ephemeral: true,
ProviderFactories: []parser.ProviderFactory{
blockingRebuildFactory{provider: provider},
},
ProviderMigrationModes: map[parser.AgentType]parser.ProviderMigrationMode{
parser.AgentCowork: parser.ProviderMigrationProviderAuthoritative,
},
},
AfterSync: func(*Engine, *db.DB) error {
firstHookCalls++
return nil
},
},
{
Name: "later",
Config: EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentCowork: {root}},
Machine: "later", Ephemeral: true,
ProviderFactories: []parser.ProviderFactory{
trackingRebuildFactory{provider: secondProvider},
},
ProviderMigrationModes: map[parser.AgentType]parser.ProviderMigrationMode{
parser.AgentCowork: parser.ProviderMigrationProviderAuthoritative,
},
},
AfterSync: func(*Engine, *db.DB) error {
secondHookCalls++
return nil
},
},
},
})
result <- struct {
stats SyncStats
err error
}{stats: stats, err: runErr}
}()
select {
case <-provider.started:
case <-time.After(5 * time.Second):
t.Fatal("contributor parse did not reach controlled boundary")
}
cancel()
select {
case got := <-result:
require.ErrorIs(t, got.err, context.Canceled)
assert.True(t, got.stats.Aborted)
assert.Zero(t, firstHookCalls, "AfterSync ran on incomplete contributor data")
assert.Zero(t, secondProvider.discoverCalls, "later contributor started after cancellation")
assert.Zero(t, secondHookCalls, "later contributor hook ran after cancellation")
case <-time.After(5 * time.Second):
t.Fatal("cancelled contributor rebuild did not return")
}
page, searchErr := database.Search(context.Background(), db.SearchFilter{
Query: "archive before cancel", Limit: 5,
})
require.NoError(t, searchErr)
require.Len(t, page.Results, 1)
assert.NoFileExists(t, database.Path()+resyncTempSuffix)
assert.NoFileExists(t, database.Path()+resyncTempSuffix+"-wal")
assert.NoFileExists(t, database.Path()+resyncTempSuffix+"-shm")
}
func TestResyncLocalCancellationPreventsContributors(t *testing.T) {
root := t.TempDir()
database, err := db.Open(filepath.Join(t.TempDir(), "archive.db"))
require.NoError(t, err)
t.Cleanup(func() { require.NoError(t, database.Close()) })
oldPath := filepath.Join(root, "old", "old.jsonl")
require.NoError(t, os.MkdirAll(filepath.Dir(oldPath), 0o755))
require.NoError(t, os.WriteFile(oldPath, []byte(testjsonl.NewSessionBuilder().
AddClaudeUser("2026-01-01T00:00:00Z", "archive before local cancel").String()), 0o644))
seedEngine := NewEngine(database, EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentClaude: {root}},
Machine: "local",
})
require.Equal(t, 1, seedEngine.SyncAll(context.Background(), nil).Synced)
seedEngine.Close()
require.NoError(t, os.Remove(oldPath))
blocking := &blockingRebuildProvider{
ProviderBase: parser.ProviderBase{
Def: parser.AgentDef{Type: parser.AgentCowork},
Caps: parser.Capabilities{Source: parser.SourceCapabilities{
DiscoverSources: parser.CapabilitySupported,
}},
},
started: make(chan struct{}),
}
later := &trackingRebuildProvider{ProviderBase: parser.ProviderBase{
Def: parser.AgentDef{Type: parser.AgentCowork},
Caps: parser.Capabilities{Source: parser.SourceCapabilities{
DiscoverSources: parser.CapabilitySupported,
}},
}}
engine := NewEngine(database, EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentCowork: {root}},
Machine: "local",
ProviderFactories: []parser.ProviderFactory{
blockingRebuildFactory{provider: blocking},
},
ProviderMigrationModes: map[parser.AgentType]parser.ProviderMigrationMode{
parser.AgentCowork: parser.ProviderMigrationProviderAuthoritative,
},
})
t.Cleanup(engine.Close)
ctx, cancel := context.WithCancel(context.Background())
result := make(chan struct {
stats SyncStats
err error
}, 1)
go func() {
stats, runErr := engine.ResyncAllWithOptions(ctx, nil, RebuildOptions{
Contributors: []RebuildContributor{{
Name: "later",
Config: EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentCowork: {root}},
Machine: "later", Ephemeral: true,
ProviderFactories: []parser.ProviderFactory{
trackingRebuildFactory{provider: later},
},
ProviderMigrationModes: map[parser.AgentType]parser.ProviderMigrationMode{
parser.AgentCowork: parser.ProviderMigrationProviderAuthoritative,
},
},
}},
})
result <- struct {
stats SyncStats
err error
}{stats: stats, err: runErr}
}()
select {
case <-blocking.started:
case <-time.After(5 * time.Second):
t.Fatal("local parse did not reach controlled boundary")
}
cancel()
select {
case got := <-result:
require.ErrorIs(t, got.err, context.Canceled)
assert.True(t, got.stats.Aborted)
assert.Zero(t, later.discoverCalls)
case <-time.After(5 * time.Second):
t.Fatal("cancelled local rebuild did not return")
}
page, searchErr := database.Search(context.Background(), db.SearchFilter{
Query: "archive before local cancel", Limit: 5,
})
require.NoError(t, searchErr)
require.Len(t, page.Results, 1)
assert.NoFileExists(t, database.Path()+resyncTempSuffix)
assert.NoFileExists(t, database.Path()+resyncTempSuffix+"-wal")
assert.NoFileExists(t, database.Path()+resyncTempSuffix+"-shm")
}
type staticUsageRebuildProvider struct {
parser.ProviderBase
source parser.SourceRef
result parser.ParseResult
}
func (p *staticUsageRebuildProvider) Discover(context.Context) ([]parser.SourceRef, error) {
return []parser.SourceRef{p.source}, nil
}
func (p *staticUsageRebuildProvider) Fingerprint(
context.Context, parser.SourceRef,
) (parser.SourceFingerprint, error) {
return parser.SourceFingerprint{Key: p.source.FingerprintKey, Size: 10, MTimeNS: 10}, nil
}
func (p *staticUsageRebuildProvider) Parse(
context.Context, parser.ParseRequest,
) (parser.ParseOutcome, error) {
return parser.ParseOutcome{
Results: []parser.ParseResultOutcome{{Result: p.result}},
ResultSetComplete: true,
}, nil
}
type staticUsageRebuildFactory struct{ provider *staticUsageRebuildProvider }
func (f staticUsageRebuildFactory) Definition() parser.AgentDef {
return f.provider.Definition()
}
func (f staticUsageRebuildFactory) Capabilities() parser.Capabilities {
return f.provider.Capabilities()
}
func (f staticUsageRebuildFactory) NewProvider(parser.ProviderConfig) parser.Provider {
return f.provider
}
func newStaticUsageRebuildProvider(id, path string, input, output int) *staticUsageRebuildProvider {
started := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC)
return &staticUsageRebuildProvider{
ProviderBase: parser.ProviderBase{
Def: parser.AgentDef{Type: parser.AgentCowork, FileBased: true},
Caps: parser.Capabilities{Source: parser.SourceCapabilities{
DiscoverSources: parser.CapabilitySupported,
CompositeFingerprint: parser.CapabilitySupported,
}},
},
source: parser.SourceRef{
Provider: parser.AgentCowork, Key: path, DisplayPath: path, FingerprintKey: path,
},
result: parser.ParseResult{
Session: parser.ParsedSession{
ID: id, Project: "usage", Machine: "fixture", Agent: parser.AgentCowork,
StartedAt: started, EndedAt: started, FirstMessage: "usage fixture",
MessageCount: 1, UserMessageCount: 1,
File: parser.FileInfo{Path: path, Size: 10, Mtime: 10},
},
Messages: []parser.ParsedMessage{{
Ordinal: 0, Role: parser.RoleUser, Content: "usage fixture", Timestamp: started,
}},
UsageEvents: []parser.ParsedUsageEvent{{
SessionID: id, Source: "fixture", Model: "fixture-model",
InputTokens: input, OutputTokens: output,
OccurredAt: "2026-01-01T00:00:00Z", DedupKey: id + "-usage",
}},
},
}
}
func TestResyncContributorBatchFailureAbortsAndCleansTempDB(t *testing.T) {
root := t.TempDir()
database, err := db.Open(filepath.Join(t.TempDir(), "archive.db"))
require.NoError(t, err)
t.Cleanup(func() { require.NoError(t, database.Close()) })
engine := NewEngine(database, EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentClaude: {root}},
Machine: "local",
})
t.Cleanup(engine.Close)
oldPath := filepath.Join(root, "old", "old.jsonl")
require.NoError(t, os.MkdirAll(filepath.Dir(oldPath), 0o755))
require.NoError(t, os.WriteFile(oldPath, []byte(testjsonl.NewSessionBuilder().
AddClaudeUser("2026-01-01T00:00:00Z", "archive before batch failure").String()), 0o644))
require.Equal(t, 1, engine.SyncAll(context.Background(), nil).Synced)
require.NoError(t, os.Remove(oldPath))
bad := newStaticUsageRebuildProvider("batch-rejected", "bad-source", 1, 1)
bad.result.UsageEvents = nil
bad.result.Messages = append(bad.result.Messages, parser.ParsedMessage{
Ordinal: 0, Role: parser.RoleAssistant, Content: "duplicate ordinal",
Timestamp: time.Date(2026, 1, 1, 0, 0, 1, 0, time.UTC),
})
stats, err := engine.ResyncAllWithOptions(context.Background(), nil, RebuildOptions{
Contributors: []RebuildContributor{{
Name: "bad-batch",
Config: EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentCowork: {"bad-root"}},
Machine: "bad-batch", Ephemeral: true,
ProviderFactories: []parser.ProviderFactory{
staticUsageRebuildFactory{provider: bad},
},
ProviderMigrationModes: map[parser.AgentType]parser.ProviderMigrationMode{
parser.AgentCowork: parser.ProviderMigrationProviderAuthoritative,
},
},
}},
})
require.NoError(t, err)
assert.True(t, stats.Aborted)
assert.Equal(t, 1, stats.Failed)
page, searchErr := database.Search(context.Background(), db.SearchFilter{
Query: "archive before batch failure", Limit: 5,
})
require.NoError(t, searchErr)
require.Len(t, page.Results, 1)
assert.NoFileExists(t, database.Path()+resyncTempSuffix)
assert.NoFileExists(t, database.Path()+resyncTempSuffix+"-wal")
assert.NoFileExists(t, database.Path()+resyncTempSuffix+"-shm")
}
type malformedRebuildProvider struct{ parser.ProviderBase }
func (p *malformedRebuildProvider) Discover(context.Context) ([]parser.SourceRef, error) {
out := make([]parser.SourceRef, 3)
for i := range out {
key := fmt.Sprintf("malformed-%d.jsonl", i)
out[i] = parser.SourceRef{
Provider: parser.AgentCowork, Key: key, DisplayPath: key, FingerprintKey: key,
}
}
return out, nil
}
func (p *malformedRebuildProvider) Fingerprint(
_ context.Context, source parser.SourceRef,
) (parser.SourceFingerprint, error) {
return parser.SourceFingerprint{Key: source.FingerprintKey, Size: 8, MTimeNS: 10}, nil
}
func (p *malformedRebuildProvider) Parse(
_ context.Context, request parser.ParseRequest,
) (parser.ParseOutcome, error) {
return parser.ParseOutcome{}, fmt.Errorf("malformed fixture %s", request.Source.Key)
}
type malformedRebuildFactory struct{ provider *malformedRebuildProvider }
func (f malformedRebuildFactory) Definition() parser.AgentDef {
return f.provider.Definition()
}
func (f malformedRebuildFactory) Capabilities() parser.Capabilities {
return f.provider.Capabilities()
}
func (f malformedRebuildFactory) NewProvider(parser.ProviderConfig) parser.Provider {
return f.provider
}
func TestResyncContributorParserFailuresAbortAndCleanTempDB(t *testing.T) {
root := t.TempDir()
database, err := db.Open(filepath.Join(t.TempDir(), "archive.db"))
require.NoError(t, err)
t.Cleanup(func() { require.NoError(t, database.Close()) })
engine := NewEngine(database, EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentClaude: {root}},
Machine: "local",
})
t.Cleanup(engine.Close)
oldPath := filepath.Join(root, "old", "old.jsonl")
require.NoError(t, os.MkdirAll(filepath.Dir(oldPath), 0o755))
require.NoError(t, os.WriteFile(oldPath, []byte(testjsonl.NewSessionBuilder().
AddClaudeUser("2026-01-01T00:00:00Z", "archive before parser failures").String()), 0o644))
require.Equal(t, 1, engine.SyncAll(context.Background(), nil).Synced)
require.NoError(t, os.Remove(oldPath))
provider := &malformedRebuildProvider{ProviderBase: parser.ProviderBase{
Def: parser.AgentDef{Type: parser.AgentCowork, FileBased: true},
Caps: parser.Capabilities{Source: parser.SourceCapabilities{
DiscoverSources: parser.CapabilitySupported,
CompositeFingerprint: parser.CapabilitySupported,
}},
}}
stats, err := engine.ResyncAllWithOptions(context.Background(), nil, RebuildOptions{
Contributors: []RebuildContributor{{
Name: "malformed",
Config: EngineConfig{
AgentDirs: map[parser.AgentType][]string{parser.AgentCowork: {"malformed-root"}},
Machine: "malformed", Ephemeral: true,
ProviderFactories: []parser.ProviderFactory{
malformedRebuildFactory{provider: provider},
},
ProviderMigrationModes: map[parser.AgentType]parser.ProviderMigrationMode{
parser.AgentCowork: parser.ProviderMigrationProviderAuthoritative,
},
},
}},
})
require.NoError(t, err)
assert.True(t, stats.Aborted)
assert.Equal(t, 3, stats.Failed)
assert.Equal(t, 3, stats.TotalSessions)
page, searchErr := database.Search(context.Background(), db.SearchFilter{
Query: "archive before parser failures", Limit: 5,
})
require.NoError(t, searchErr)
require.Len(t, page.Results, 1)
assert.NoFileExists(t, database.Path()+resyncTempSuffix)
assert.NoFileExists(t, database.Path()+resyncTempSuffix+"-wal")
assert.NoFileExists(t, database.Path()+resyncTempSuffix+"-shm")
}