diff --git a/cmd/agentsview/archive_query_backend.go b/cmd/agentsview/archive_query_backend.go index b3c96c33d..b38a7dc21 100644 --- a/cmd/agentsview/archive_query_backend.go +++ b/cmd/agentsview/archive_query_backend.go @@ -249,6 +249,7 @@ func (b localArchiveQueryBackend) SessionUsage( if known && !b.skipFreshData { engine := sync.NewEngine(b.database, sync.EngineConfig{ AgentDirs: b.cfg.AgentDirs, + IncludeCwdPrefixes: b.cfg.SyncIncludeCwdPrefixes, Machine: "local", BlockedResultCategories: b.cfg.ResultContentBlockedCategories, }) diff --git a/cmd/agentsview/archive_write_backend.go b/cmd/agentsview/archive_write_backend.go index 34948d165..fa2fb6895 100644 --- a/cmd/agentsview/archive_write_backend.go +++ b/cmd/agentsview/archive_write_backend.go @@ -549,6 +549,7 @@ func (b *localArchiveWriteBackend) PGPushWatch( engine := syncpkg.NewEngine(b.database, syncpkg.EngineConfig{ AgentDirs: b.appCfg.AgentDirs, + IncludeCwdPrefixes: b.appCfg.SyncIncludeCwdPrefixes, Machine: "local", BlockedResultCategories: b.appCfg.ResultContentBlockedCategories, }) diff --git a/cmd/agentsview/cli.go b/cmd/agentsview/cli.go index 7884e999b..adf9611a1 100644 --- a/cmd/agentsview/cli.go +++ b/cmd/agentsview/cli.go @@ -850,6 +850,13 @@ func writeRootHelp(w io.Writer, root *cobra.Command) { fmt.Fprintln(w, " Example:") fmt.Fprintln(w, " watch_exclude_patterns = [\".git\", \"node_modules\", \".next\", \"dist\"]") fmt.Fprintln(w) + fmt.Fprintln(w, "Session cwd filter:") + fmt.Fprintln(w, " Add \"sync_include_cwd_prefixes\" to ~/.agentsview/config.toml to") + fmt.Fprintln(w, " ingest only sessions whose working directory is under one of the") + fmt.Fprintln(w, " listed paths. Sessions without a recorded cwd are skipped while") + fmt.Fprintln(w, " the filter is set. Applies to local sync only. Example:") + fmt.Fprintln(w, " sync_include_cwd_prefixes = [\"/home/me/work\"]") + fmt.Fprintln(w) fmt.Fprintln(w, "Multiple directories:") fmt.Fprintln(w, " Add arrays to ~/.agentsview/config.toml to scan multiple locations:") fmt.Fprintln(w, " claude_project_dirs = [\"/path/one\", \"/path/two\"]") diff --git a/cmd/agentsview/main.go b/cmd/agentsview/main.go index d56ec5503..b2de32a03 100644 --- a/cmd/agentsview/main.go +++ b/cmd/agentsview/main.go @@ -207,6 +207,7 @@ func runServe(cfg config.Config, opts serveOptions) { if !cfg.NoSync { engine = sync.NewEngine(database, sync.EngineConfig{ AgentDirs: cfg.AgentDirs, + IncludeCwdPrefixes: cfg.SyncIncludeCwdPrefixes, Machine: "local", BlockedResultCategories: cfg.ResultContentBlockedCategories, Emitter: emitter, diff --git a/cmd/agentsview/parse_diff.go b/cmd/agentsview/parse_diff.go index a6a92bcc4..2abadaf00 100644 --- a/cmd/agentsview/parse_diff.go +++ b/cmd/agentsview/parse_diff.go @@ -127,6 +127,7 @@ func doParseDiff(cfg ParseDiffConfig) (failed bool) { engine := sync.NewDiffEngine(database, sync.EngineConfig{ AgentDirs: appCfg.AgentDirs, + IncludeCwdPrefixes: appCfg.SyncIncludeCwdPrefixes, Machine: "local", BlockedResultCategories: appCfg.ResultContentBlockedCategories, }) diff --git a/cmd/agentsview/session_sync.go b/cmd/agentsview/session_sync.go index 314cdd96d..75aa8ad92 100644 --- a/cmd/agentsview/session_sync.go +++ b/cmd/agentsview/session_sync.go @@ -73,8 +73,9 @@ func syncService( return nil, nil, fmt.Errorf("opening db: %w", err) } engine := sync.NewEngine(d, sync.EngineConfig{ - AgentDirs: cfg.AgentDirs, - Machine: "local", + AgentDirs: cfg.AgentDirs, + IncludeCwdPrefixes: cfg.SyncIncludeCwdPrefixes, + Machine: "local", }) // Close the engine before the DB so pending debounced signal // recomputes flush while the DB is still open. diff --git a/cmd/agentsview/sync.go b/cmd/agentsview/sync.go index 70f8aea62..983bd71c8 100644 --- a/cmd/agentsview/sync.go +++ b/cmd/agentsview/sync.go @@ -464,6 +464,7 @@ func runLocalSync( engine := sync.NewEngine(database, sync.EngineConfig{ AgentDirs: appCfg.AgentDirs, + IncludeCwdPrefixes: appCfg.SyncIncludeCwdPrefixes, Machine: "local", BlockedResultCategories: appCfg.ResultContentBlockedCategories, }) diff --git a/cmd/agentsview/usage.go b/cmd/agentsview/usage.go index c8b605ae8..2d9c3fc68 100644 --- a/cmd/agentsview/usage.go +++ b/cmd/agentsview/usage.go @@ -309,8 +309,9 @@ func ensureFreshData( if database.NeedsResync() { engine := sync.NewEngine(database, sync.EngineConfig{ - AgentDirs: appCfg.AgentDirs, - Machine: "local", + AgentDirs: appCfg.AgentDirs, + IncludeCwdPrefixes: appCfg.SyncIncludeCwdPrefixes, + Machine: "local", }) defer engine.Close() fmt.Fprintln(os.Stderr, @@ -332,8 +333,9 @@ func ensureFreshData( } engine := sync.NewEngine(database, sync.EngineConfig{ - AgentDirs: appCfg.AgentDirs, - Machine: "local", + AgentDirs: appCfg.AgentDirs, + IncludeCwdPrefixes: appCfg.SyncIncludeCwdPrefixes, + Machine: "local", }) defer engine.Close() diff --git a/docs/configuration.md b/docs/configuration.md index f94c381dd..cf7f19846 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -681,6 +681,42 @@ state while they hold the start lock, so `agentsview serve status` can show the starting PID, elapsed time, current phase, progress detail, and log path before the HTTP server is ready. +### Restricting Ingestion by Working Directory + +By default every discovered session is ingested. To limit the archive to +sessions from specific workspaces — for example on a machine shared across +multiple clients where transcripts from one workspace should never appear +alongside another — set `sync_include_cwd_prefixes` in +`~/.agentsview/config.toml`: + +```toml +sync_include_cwd_prefixes = [ + "/home/me/work/client-a", + "/home/me/oss", +] +``` + +When the list is non-empty, a session is ingested only if its recorded +working directory equals one of the prefixes or lives underneath one. +Prefixes and session directories are lexically cleaned before matching: +trailing separators are ignored and `..` components are resolved, so +`/home/me/oss/../other` does not match a `/home/me/oss` prefix. Matching +is path-boundary aware (`/home/me/oss` matches `/home/me/oss/repo` but +not `/home/me/oss-other`), case-sensitive, and uses the local operating +system's path separator — on Linux and macOS a backslash is an ordinary +filename character, not a directory boundary. Use absolute paths; `~` is +not expanded. + +Notes: + +- Sessions without a recorded working directory (a few agents do not store + one) are skipped while the filter is set. +- The filter gates ingestion only. Sessions already in the archive are + preserved (the SQLite database is a persistent archive); remove unwanted + existing sessions explicitly with `agentsview prune`. +- Remote-host sync is unaffected: the prefixes describe local paths, so they + are not applied to sessions pulled from `[[remote_hosts]]` entries. + ### Large Watch Trees The recursive watcher has a hard budget of 8192 directories per process. If a diff --git a/internal/config/config.go b/internal/config/config.go index 91cd63d9b..7ebac9854 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -451,6 +451,14 @@ type Config struct { ResultContentBlockedCategories []string `json:"result_content_blocked_categories,omitempty" toml:"result_content_blocked_categories"` + // SyncIncludeCwdPrefixes, when non-empty, restricts local session + // ingestion to sessions whose working directory equals one of the + // prefixes or lives underneath one. Sessions without a recorded + // cwd are skipped while the filter is active. Config-file only; + // remote sync is unaffected because the prefixes describe local + // paths. + SyncIncludeCwdPrefixes []string `json:"-" toml:"sync_include_cwd_prefixes"` + // EventsCoalesceInterval is the minimum wall-clock time between // SSE data_changed broadcasts to connected clients. Emits that // arrive within this window after a prior broadcast are coalesced @@ -944,6 +952,7 @@ func (c *Config) applyConfigTOML(data string) error { PublicOrigins []string `toml:"public_origins"` Proxy ProxyConfig `toml:"proxy"` WatchExcludePatterns []string `toml:"watch_exclude_patterns"` + SyncIncludeCwdPrefixes []string `toml:"sync_include_cwd_prefixes"` ResultContentBlockedCategories []string `toml:"result_content_blocked_categories"` Terminal TerminalConfig `toml:"terminal"` AuthToken string `toml:"auth_token"` @@ -1005,6 +1014,9 @@ func (c *Config) applyConfigTOML(data string) error { if file.WatchExcludePatterns != nil { c.WatchExcludePatterns = file.WatchExcludePatterns } + if file.SyncIncludeCwdPrefixes != nil { + c.SyncIncludeCwdPrefixes = file.SyncIncludeCwdPrefixes + } if file.ResultContentBlockedCategories != nil { c.ResultContentBlockedCategories = file.ResultContentBlockedCategories } diff --git a/internal/config/config_test.go b/internal/config/config_test.go index df3f38f9e..ba40129c7 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -1733,3 +1733,26 @@ func TestValidateRemoteHosts(t *testing.T) { }) } } + +func TestLoadFile_SyncIncludeCwdPrefixes(t *testing.T) { + f := newConfigFixture(t) + f.WriteConfigText(t, `sync_include_cwd_prefixes = ["/home/me/work", "/home/me/oss"] +`) + + cfg := f.LoadMinimal(t) + + assert.Equal(t, + []string{"/home/me/work", "/home/me/oss"}, + cfg.SyncIncludeCwdPrefixes, + ) +} + +func TestLoadFile_SyncIncludeCwdPrefixesDefaultsEmpty(t *testing.T) { + f := newConfigFixture(t) + f.WriteConfigText(t, `host = "127.0.0.1" +`) + + cfg := f.LoadMinimal(t) + + assert.Empty(t, cfg.SyncIncludeCwdPrefixes) +} diff --git a/internal/server/huma_routes_sync.go b/internal/server/huma_routes_sync.go index 58efceb6b..7eac26a0a 100644 --- a/internal/server/huma_routes_sync.go +++ b/internal/server/huma_routes_sync.go @@ -156,6 +156,7 @@ func (s *Server) syncEngineForLocal(local *db.DB) *syncpkg.Engine { } s.onDemandEngine = syncpkg.NewEngine(local, syncpkg.EngineConfig{ AgentDirs: s.cfg.AgentDirs, + IncludeCwdPrefixes: s.cfg.SyncIncludeCwdPrefixes, Machine: "local", BlockedResultCategories: s.cfg.ResultContentBlockedCategories, Emitter: emitter, diff --git a/internal/sync/cwd_filter.go b/internal/sync/cwd_filter.go new file mode 100644 index 000000000..abd783c1c --- /dev/null +++ b/internal/sync/cwd_filter.go @@ -0,0 +1,114 @@ +package sync + +import ( + "os" + "path/filepath" + "strings" + + "go.kenn.io/agentsview/internal/parser" +) + +// cwdPrefixFilter gates session ingestion on the session working +// directory. An empty filter allows everything. A non-empty filter +// allows only sessions whose cwd equals a configured prefix or lives +// underneath one; sessions without a recorded cwd are rejected +// because they cannot be attributed to any workspace. +// +// Prefixes and cwds are lexically cleaned before matching and the +// path boundary is the local OS separator. The filter only ever sees +// local paths (remote sync does not apply it), so local filesystem +// semantics are the correct ones: on POSIX a backslash is an ordinary +// filename character, not a boundary, and a cwd like "/a/b/../c" is +// resolved to "/a/c" rather than prefix-matching "/a/b". +type cwdPrefixFilter struct { + prefixes []string +} + +// newCwdPrefixFilter normalizes the configured prefixes: entries are +// trimmed, blank entries are dropped, and each remaining entry is +// cleaned with filepath.Clean so "/a/b/" and "/a/b" behave +// identically and ".." components cannot linger in a prefix. +func newCwdPrefixFilter(prefixes []string) cwdPrefixFilter { + normalized := make([]string, 0, len(prefixes)) + for _, p := range prefixes { + p = strings.TrimSpace(p) + if p == "" { + continue + } + normalized = append(normalized, filepath.Clean(p)) + } + return cwdPrefixFilter{prefixes: normalized} +} + +func (f cwdPrefixFilter) empty() bool { + return len(f.prefixes) == 0 +} + +// allows reports whether a session with the given cwd may be +// ingested. Matching is path-boundary aware: prefix "/a/b" matches +// "/a/b" and "/a/b/c" but not "/a/bc". +func (f cwdPrefixFilter) allows(cwd string) bool { + if f.empty() { + return true + } + if cwd == "" { + return false + } + cwd = filepath.Clean(cwd) + sep := string(os.PathSeparator) + for _, p := range f.prefixes { + if cwd == p { + return true + } + prefix := p + if !strings.HasSuffix(prefix, sep) { + prefix += sep + } + if strings.HasPrefix(cwd, prefix) { + return true + } + } + return false +} + +// sourceAllowsParserExclusions reports whether a source's parser +// exclusions (including engine stale-row cleanup) may delete archived +// rows. When the cwd allow-list is active, a source proves it is +// inside the list by producing at least one allowed session or +// incremental update; a source with no allowed output is frozen — its +// exclusions would erase archived sessions whose replacement writes +// the filter vetoes, which the ingestion-only contract forbids. +// Zero-result exclusion carriers (e.g. a file that parses to no live +// session) have no cwd to judge, so they are frozen too. +func (e *Engine) sourceAllowsParserExclusions(res processResult) bool { + if e.cwdFilter.empty() { + return true + } + if res.incremental != nil && e.cwdFilter.allows(res.incremental.cwd) { + return true + } + for _, pr := range res.results { + if e.cwdFilter.allows(pr.Session.Cwd) { + return true + } + } + return false +} + +// splitResultsByCwdFilter returns the parsed sessions the cwd +// allow-list admits and the number it vetoes. With no filter +// configured it returns the input untouched. +func (e *Engine) splitResultsByCwdFilter( + results []parser.ParseResult, +) ([]parser.ParseResult, int) { + if e.cwdFilter.empty() || len(results) == 0 { + return results, 0 + } + allowed := make([]parser.ParseResult, 0, len(results)) + for _, pr := range results { + if e.cwdFilter.allows(pr.Session.Cwd) { + allowed = append(allowed, pr) + } + } + return allowed, len(results) - len(allowed) +} diff --git a/internal/sync/cwd_filter_integration_test.go b/internal/sync/cwd_filter_integration_test.go new file mode 100644 index 000000000..f306b095c --- /dev/null +++ b/internal/sync/cwd_filter_integration_test.go @@ -0,0 +1,191 @@ +package sync_test + +import ( + "context" + "os" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "go.kenn.io/agentsview/internal/dbtest" + "go.kenn.io/agentsview/internal/parser" + "go.kenn.io/agentsview/internal/sync" + "go.kenn.io/agentsview/internal/testjsonl" +) + +func setupClaudeEnvWithCwdPrefixes( + t *testing.T, prefixes []string, +) *testEnv { + t.Helper() + if testing.Short() { + t.Skip("skipping integration test") + } + env := &testEnv{db: dbtest.OpenTestDB(t), claudeDir: t.TempDir()} + env.engine = sync.NewEngine(env.db, sync.EngineConfig{ + AgentDirs: map[parser.AgentType][]string{ + parser.AgentClaude: {env.claudeDir}, + }, + Machine: "local", + IncludeCwdPrefixes: prefixes, + }) + return env +} + +func TestSyncEngineCwdPrefixFilter(t *testing.T) { + env := setupClaudeEnvWithCwdPrefixes( + t, []string{"/Users/alice/work"}, + ) + + inside := testjsonl.NewSessionBuilder(). + AddClaudeUser(tsEarly, "Inside", "/Users/alice/work/my-app"). + AddClaudeAssistant(tsEarlyS5, "ok"). + String() + outside := testjsonl.NewSessionBuilder(). + AddClaudeUser(tsEarly, "Outside", "/Users/alice/personal/blog"). + AddClaudeAssistant(tsEarlyS5, "ok"). + String() + sibling := testjsonl.NewSessionBuilder(). + AddClaudeUser(tsEarly, "Sibling", "/Users/alice/workspace"). + AddClaudeAssistant(tsEarlyS5, "ok"). + String() + + env.writeClaudeSessionForProject( + t, "/Users/alice/work/my-app", + "inside-session.jsonl", inside, + ) + env.writeClaudeSessionForProject( + t, "/Users/alice/personal/blog", + "outside-session.jsonl", outside, + ) + env.writeClaudeSessionForProject( + t, "/Users/alice/workspace", + "sibling-session.jsonl", sibling, + ) + + env.engine.SyncAll(context.Background(), nil) + + assertSessionProject(t, env.db, "inside-session", "my_app") + for _, id := range []string{"outside-session", "sibling-session"} { + sess, err := env.db.GetSession(context.Background(), id) + require.NoError(t, err, "GetSession(%q)", id) + assert.Nil(t, sess, + "session %q outside the cwd allow-list must not be ingested", id) + } +} + +// A session archived before the cwd allow-list was configured must not +// keep receiving appended messages through the incremental JSONL path, +// which bypasses the prepareSessionWrite veto. +func TestSyncEngineCwdPrefixFilterBlocksIncrementalAppend(t *testing.T) { + if testing.Short() { + t.Skip("skipping integration test") + } + + // Archive the outside-prefix session with no filter configured, + // as if it was ingested before sync_include_cwd_prefixes was set. + env := &testEnv{db: dbtest.OpenTestDB(t), claudeDir: t.TempDir()} + env.engine = sync.NewEngine(env.db, sync.EngineConfig{ + AgentDirs: map[parser.AgentType][]string{ + parser.AgentClaude: {env.claudeDir}, + }, + Machine: "local", + }) + + initial := testjsonl.NewSessionBuilder(). + AddClaudeUser(tsEarly, "Outside", "/Users/alice/personal/blog"). + AddClaudeAssistant(tsEarlyS5, "ok"). + String() + path := env.writeClaudeSessionForProject( + t, "/Users/alice/personal/blog", + "outside-append.jsonl", initial, + ) + env.engine.SyncAll(context.Background(), nil) + assertSessionMessageCount(t, env.db, "outside-append", 2) + + // Turn the filter on and append to the archived session's file. + env.engine = sync.NewEngine(env.db, sync.EngineConfig{ + AgentDirs: map[parser.AgentType][]string{ + parser.AgentClaude: {env.claudeDir}, + }, + Machine: "local", + IncludeCwdPrefixes: []string{"/Users/alice/work"}, + }) + + appended := testjsonl.ClaudeUserJSON("appended", tsEarlyS1) + "\n" + f, err := os.OpenFile(path, os.O_APPEND|os.O_WRONLY, 0o644) + require.NoError(t, err, "open for append") + _, err = f.WriteString(appended) + f.Close() + require.NoError(t, err, "append") + + env.engine.SyncPaths([]string{path}) + + // Neither the incremental path nor the full-parse fallback may + // store the appended message; the archived rows stay untouched. + assertSessionMessageCount(t, env.db, "outside-append", 2) + assertMessageRoles(t, env.db, "outside-append", "user", "assistant") +} + +// A full resync where the cwd allow-list vetoes every discovered +// session is an intentional result, not a broken rebuild: the swap +// must proceed and the orphan copy must restore the archived rows +// (the filter gates ingestion only). Without a distinct filtered +// counter the abort guard reads such a run as an unsafe empty +// rebuild and leaves NeedsResync true forever. +func TestResyncAllProceedsWhenAllSessionsCwdFiltered(t *testing.T) { + if testing.Short() { + t.Skip("skipping integration test") + } + + // Archive two sessions with no filter configured. + env := &testEnv{db: dbtest.OpenTestDB(t), claudeDir: t.TempDir()} + env.engine = sync.NewEngine(env.db, sync.EngineConfig{ + AgentDirs: map[parser.AgentType][]string{ + parser.AgentClaude: {env.claudeDir}, + }, + Machine: "local", + }) + + first := testjsonl.NewSessionBuilder(). + AddClaudeUser(tsEarly, "First", "/Users/alice/personal/blog"). + AddClaudeAssistant(tsEarlyS5, "ok"). + String() + second := testjsonl.NewSessionBuilder(). + AddClaudeUser(tsEarly, "Second", "/Users/alice/personal/notes"). + AddClaudeAssistant(tsEarlyS5, "ok"). + String() + env.writeClaudeSessionForProject( + t, "/Users/alice/personal/blog", + "filtered-one.jsonl", first, + ) + env.writeClaudeSessionForProject( + t, "/Users/alice/personal/notes", + "filtered-two.jsonl", second, + ) + env.engine.SyncAll(context.Background(), nil) + assertSessionMessageCount(t, env.db, "filtered-one", 2) + assertSessionMessageCount(t, env.db, "filtered-two", 2) + + // Resync with an allow-list that excludes every session. + env.engine = sync.NewEngine(env.db, sync.EngineConfig{ + AgentDirs: map[parser.AgentType][]string{ + parser.AgentClaude: {env.claudeDir}, + }, + Machine: "local", + IncludeCwdPrefixes: []string{"/Users/alice/work"}, + }) + stats := env.engine.ResyncAll(context.Background(), nil) + + require.False(t, stats.Aborted, + "all-filtered resync must not abort: %+v", stats.Warnings) + assert.Equal(t, 0, stats.Synced, "synced") + assert.Equal(t, 0, stats.Failed, "failed") + assert.Equal(t, 2, stats.OrphanedCopied, "orphaned copied") + + // The archived sessions survive the swap via the orphan copy. + assertSessionMessageCount(t, env.db, "filtered-one", 2) + assertSessionMessageCount(t, env.db, "filtered-two", 2) + assert.False(t, env.db.NeedsResync(), + "completed resync must clear the needs-resync marker") +} diff --git a/internal/sync/cwd_filter_test.go b/internal/sync/cwd_filter_test.go new file mode 100644 index 000000000..218ff4f98 --- /dev/null +++ b/internal/sync/cwd_filter_test.go @@ -0,0 +1,251 @@ +package sync + +import ( + "context" + "runtime" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "go.kenn.io/agentsview/internal/db" + "go.kenn.io/agentsview/internal/parser" +) + +type cwdFilterCase struct { + name string + prefixes []string + cwd string + want bool +} + +func TestCwdPrefixFilterAllows(t *testing.T) { + tests := []cwdFilterCase{ + {"empty filter allows anything", nil, "/anywhere", true}, + {"empty filter allows empty cwd", nil, "", true}, + {"exact match", []string{"/a/b"}, "/a/b", true}, + {"child path", []string{"/a/b"}, "/a/b/c/d", true}, + {"sibling with shared prefix", []string{"/a/b"}, "/a/bc", false}, + {"outside prefix", []string{"/a/b"}, "/x", false}, + {"empty cwd rejected when filter set", []string{"/a/b"}, "", false}, + {"second prefix matches", []string{"/a/b", "/x/y"}, "/x/y/z", true}, + {"trailing separator normalized", []string{"/a/b/"}, "/a/b/c", true}, + {"prefix longer than cwd", []string{"/a/b/c"}, "/a/b", false}, + {"case sensitive", []string{"/a/B"}, "/a/b/c", false}, + {"blank entries ignored", []string{" ", ""}, "/anywhere", true}, + {"root prefix allows any cwd", []string{"/"}, "/anywhere", true}, + {"dot-dot escaping the prefix rejected", []string{"/a/b"}, "/a/b/../c", false}, + {"dot-dot staying inside allowed", []string{"/a/b"}, "/a/b/c/../d", true}, + {"dot-dot in prefix cleaned", []string{"/a/b/../c"}, "/a/c/d", true}, + } + if runtime.GOOS == "windows" { + tests = append(tests, + cwdFilterCase{"backslash boundary", []string{`C:\work`}, `C:\work\repo`, true}, + cwdFilterCase{"drive sibling", []string{`C:\work`}, `C:\workspace`, false}, + cwdFilterCase{"mixed separators normalized", []string{`C:/work`}, `C:\work\repo`, true}, + ) + } else { + // On POSIX a backslash is an ordinary filename character: + // "b\evil" is a sibling of "b" under /a, not a child of /a/b. + tests = append(tests, cwdFilterCase{ + "backslash is not a separator", []string{"/a/b"}, `/a/b\evil`, false, + }) + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + f := newCwdPrefixFilter(tt.prefixes) + assert.Equal(t, tt.want, f.allows(tt.cwd)) + }) + } +} + +// exclusionGateJob builds a syncJob whose parse supersedes the +// archived "stale" row (via excludedSessionIDs) with a replacement +// session recorded at the given cwd. +func exclusionGateJob(cwd string) syncJob { + return syncJob{ + path: "/src/session.jsonl", + processResult: processResult{ + excludedSessionIDs: []string{"stale"}, + results: []parser.ParseResult{ + {Session: parser.ParsedSession{ + ID: "replacement", + Agent: parser.AgentClaude, + Machine: "local", + Project: "proj", + Cwd: cwd, + }}, + }, + }, + } +} + +// A parse whose sessions are all outside the cwd allow-list must not +// delete the archived rows its exclusion list supersedes: the +// replacement write is vetoed, so the delete would erase a session +// the filter promises to preserve. +func TestCollectAndBatchGatesParserExclusionsByCwdFilter(t *testing.T) { + ctx := context.Background() + + t.Run("filtered source keeps archived row", func(t *testing.T) { + database := openTestDB(t) + require.NoError(t, database.UpsertSession(db.Session{ + ID: "stale", Project: "proj", Machine: "local", Agent: "claude", + })) + e := NewEngine(database, EngineConfig{ + Machine: "local", + IncludeCwdPrefixes: []string{"/allowed"}, + }) + + results := make(chan syncJob, 1) + results <- exclusionGateJob("/outside/repo") + close(results) + stats := e.collectAndBatch( + ctx, results, 1, 1, nil, syncWriteDefault, + ) + + gotStale, err := database.GetSession(ctx, "stale") + require.NoError(t, err) + assert.NotNil(t, gotStale, + "archived row must survive exclusions from a filtered source") + gotNew, err := database.GetSession(ctx, "replacement") + require.NoError(t, err) + assert.Nil(t, gotNew, "filtered replacement must not be written") + assert.Empty(t, stats.parserExcludedIDs, + "frozen exclusions must not reach resync orphan-copy exclusion") + assert.Equal(t, 1, stats.cwdFilteredSessions, "filtered sessions") + assert.Equal(t, 1, stats.cwdFilteredFiles, "filtered files") + assert.Equal(t, 0, stats.Synced, "synced") + }) + + t.Run("allowed source deletes superseded row", func(t *testing.T) { + database := openTestDB(t) + require.NoError(t, database.UpsertSession(db.Session{ + ID: "stale", Project: "proj", Machine: "local", Agent: "claude", + })) + e := NewEngine(database, EngineConfig{ + Machine: "local", + IncludeCwdPrefixes: []string{"/allowed"}, + }) + + results := make(chan syncJob, 1) + results <- exclusionGateJob("/allowed/repo") + close(results) + stats := e.collectAndBatch( + ctx, results, 1, 1, nil, syncWriteDefault, + ) + + gotStale, err := database.GetSession(ctx, "stale") + require.NoError(t, err) + assert.Nil(t, gotStale, + "superseded row must be deleted for an allowed source") + gotNew, err := database.GetSession(ctx, "replacement") + require.NoError(t, err) + assert.NotNil(t, gotNew, "allowed replacement must be written") + assert.Equal(t, []string{"stale"}, stats.parserExcludedIDs) + assert.Equal(t, 0, stats.cwdFilteredSessions, "filtered sessions") + assert.Equal(t, 1, stats.Synced, "synced") + }) +} + +func TestShouldAbortResyncSwap(t *testing.T) { + tests := []struct { + name string + stats SyncStats + oldFileSessions int + trashedCopied int + want bool + }{ + { + name: "clean run proceeds", + stats: SyncStats{ + TotalSessions: 5, Synced: 5, + filesOK: 5, nonContainerDiscovered: 5, + }, + oldFileSessions: 5, + }, + { + name: "cancelled run aborts", + stats: SyncStats{Aborted: true, Synced: 5}, + oldFileSessions: 5, + want: true, + }, + { + name: "empty discovery with old data aborts", + stats: SyncStats{}, + oldFileSessions: 3, + want: true, + }, + { + name: "zero writes unexplained aborts", + stats: SyncStats{ + TotalSessions: 3, nonContainerDiscovered: 3, + }, + oldFileSessions: 3, + want: true, + }, + { + name: "more failures than successes aborts", + stats: SyncStats{ + TotalSessions: 6, Synced: 1, Failed: 5, + filesOK: 1, nonContainerDiscovered: 6, + }, + oldFileSessions: 6, + want: true, + }, + { + name: "parser-excluded-only run proceeds", + stats: SyncStats{ + TotalSessions: 3, filesOK: 3, + parserExcludedFiles: 3, nonContainerDiscovered: 3, + }, + oldFileSessions: 3, + }, + { + name: "all-cwd-filtered run proceeds", + stats: SyncStats{ + TotalSessions: 2, filesOK: 2, + cwdFilteredFiles: 2, cwdFilteredSessions: 2, + nonContainerDiscovered: 2, + }, + oldFileSessions: 2, + }, + { + name: "cwd-filtered mixed with parser-excluded proceeds", + stats: SyncStats{ + TotalSessions: 4, filesOK: 4, + cwdFilteredFiles: 2, cwdFilteredSessions: 3, + parserExcludedFiles: 2, nonContainerDiscovered: 4, + }, + oldFileSessions: 4, + }, + { + name: "cwd-filtered with unaccounted OK file aborts", + stats: SyncStats{ + TotalSessions: 3, filesOK: 2, + cwdFilteredFiles: 1, cwdFilteredSessions: 1, + nonContainerDiscovered: 3, + }, + oldFileSessions: 3, + want: true, + }, + { + name: "cwd-filtered with failures aborts", + stats: SyncStats{ + TotalSessions: 2, Failed: 1, filesOK: 1, + cwdFilteredFiles: 1, cwdFilteredSessions: 1, + nonContainerDiscovered: 2, + }, + oldFileSessions: 2, + want: true, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got := shouldAbortResyncSwap( + tt.stats, tt.oldFileSessions, tt.trashedCopied, + ) + assert.Equal(t, tt.want, got) + }) + } +} diff --git a/internal/sync/engine.go b/internal/sync/engine.go index c6d53c2c8..67c6105e5 100644 --- a/internal/sync/engine.go +++ b/internal/sync/engine.go @@ -60,6 +60,14 @@ type EngineConfig struct { AgentDirs map[parser.AgentType][]string Machine string BlockedResultCategories []string + // IncludeCwdPrefixes, when non-empty, restricts ingestion to + // sessions whose working directory equals one of the prefixes + // or lives underneath one. Sessions without a recorded cwd are + // skipped while the filter is active. Populated from the + // sync_include_cwd_prefixes config option for local sync; + // remote sync leaves it empty because the prefixes describe + // local paths. + IncludeCwdPrefixes []string // IDPrefix is prepended to all session IDs. Used by // remote sync to namespace IDs by host (e.g. "host~"). IDPrefix string @@ -90,6 +98,7 @@ type Engine struct { agentDirs map[parser.AgentType][]string machine string blockedResultCategories map[string]bool + cwdFilter cwdPrefixFilter syncMu gosync.Mutex // serializes all sync operations mu gosync.RWMutex lastSync time.Time @@ -222,6 +231,7 @@ func NewEngine( agentDirs: dirs, machine: cfg.Machine, blockedResultCategories: blockedCategorySet(cfg.BlockedResultCategories), + cwdFilter: newCwdPrefixFilter(cfg.IncludeCwdPrefixes), skipCache: skipCache, skipFingerprints: make(map[string]string), s3CodexIndexCache: make(map[string]s3CodexIndexSnapshot), @@ -1031,6 +1041,61 @@ const resyncTempSuffix = "-resync" // atomically swaps the files and reopens the original DB // handle. This avoids the per-row trigger overhead of bulk // deleting hundreds of thousands of messages in place. +// shouldAbortResyncSwap decides whether a finished resync pass built a +// database that would be worse than the original, so the swap must be +// abandoned: +// - sync was cancelled (partial rebuild) +// - nothing synced at all (empty discovery, or all skipped) +// when old DB had data +// - more files failed than succeeded (permission errors, +// disk issues) +// +// OpenCode-only rebuilds are allowed to finish with 0 freshly synced +// sessions when every storage parse was intentionally preserved +// against the archive; orphan copy restores those rows immediately +// after the sync pass. A few permanent parse failures are tolerated +// since those files were broken in the old DB too. +// OpenCode-format storage is a self-preserving container store that +// flows through file discovery, so it is excluded from the discovery +// check just as it is subtracted from oldFileSessions by the caller. +// Otherwise its discovery would mask the disappearance of plain +// file-backed sessions whose directories went empty. +func shouldAbortResyncSwap( + stats SyncStats, oldFileSessions, trashedCopied int, +) bool { + emptyDiscovery := stats.nonContainerDiscovered == 0 && + oldFileSessions > 0 + preservedOnly := stats.Synced == 0 && + stats.TotalSessions > 0 && + stats.Failed == 0 && + (oldFileSessions == 0 || trashedCopied > 0) + excludedOnly := stats.Synced == 0 && + stats.TotalSessions > 0 && + stats.Failed == 0 && + stats.parserExcludedFiles > 0 && + stats.filesOK == stats.parserExcludedFiles + // A zero-write run is intentional when the sync_include_cwd_prefixes + // allow-list vetoed sessions AND every OK file is accounted for as + // either fully filtered or parser-excluded: the swap proceeds and + // the orphan copy restores the archived rows, because the filter + // gates ingestion only. Requiring the full accounting keeps the + // guard armed for mixed runs where other files produced nothing for + // an unexplained reason. + cwdFilteredOnly := stats.Synced == 0 && + stats.TotalSessions > 0 && + stats.Failed == 0 && + stats.cwdFilteredSessions > 0 && + stats.filesOK == stats.cwdFilteredFiles+stats.parserExcludedFiles + return stats.Aborted || + emptyDiscovery || + (stats.Synced == 0 && + stats.TotalSessions > 0 && + !preservedOnly && + !excludedOnly && + !cwdFilteredOnly) || + (stats.Failed > 0 && stats.Failed > stats.filesOK) +} + func (e *Engine) ResyncAll( ctx context.Context, onProgress ProgressFunc, ) (stats SyncStats) { @@ -1221,42 +1286,7 @@ func (e *Engine) resyncAllLocked( e.openCodeArchiveStore = nil e.phaseStats.Log("resync") - // Abort swap when the fresh DB would be worse than the - // original: - // - sync was cancelled (partial rebuild) - // - nothing synced at all (empty discovery, or all skipped) - // when old DB had data - // - more files failed than succeeded (permission errors, - // disk issues) - // OpenCode-only rebuilds are allowed to finish with 0 - // freshly synced sessions when every storage parse was - // intentionally preserved against the archive; orphan copy - // restores those rows immediately after the sync pass. - // A few permanent parse failures are tolerated since those - // files were broken in the old DB too. - // OpenCode-format storage is a self-preserving container store that - // now flows through file discovery, so it is excluded here just as it - // is subtracted from oldFileSessions above. Otherwise its discovery - // would mask the disappearance of plain file-backed sessions whose - // directories went empty. - emptyDiscovery := stats.nonContainerDiscovered == 0 && - oldFileSessions > 0 - preservedOnly := stats.Synced == 0 && - stats.TotalSessions > 0 && - stats.Failed == 0 && - (oldFileSessions == 0 || trashedCopied > 0) - excludedOnly := stats.Synced == 0 && - stats.TotalSessions > 0 && - stats.Failed == 0 && - stats.parserExcludedFiles > 0 && - stats.filesOK == stats.parserExcludedFiles - abortSwap := stats.Aborted || - emptyDiscovery || - (stats.Synced == 0 && - stats.TotalSessions > 0 && - !preservedOnly && - !excludedOnly) || - (stats.Failed > 0 && stats.Failed > stats.filesOK) + abortSwap := shouldAbortResyncSwap(stats, oldFileSessions, trashedCopied) if abortSwap { log.Printf( "resync: aborting swap, %d synced / %d failed / %d total", @@ -3134,13 +3164,14 @@ func (e *Engine) syncProviderDBBackedAgent( tWrite := time.Now() var written int if writeMode == syncWriteBulk { - var failedWrites int - written, _, failedWrites = e.writeBatch( + var failedWrites, cwdFiltered int + written, _, failedWrites, cwdFiltered = e.writeBatch( pending, writeMode, true, ) for range failedWrites { stats.RecordFailed() } + stats.cwdFilteredSessions += cwdFiltered } else { resolveWorktreeProject := e.loadWorktreeProjectResolver() for _, pw := range pending { @@ -3293,6 +3324,16 @@ func (e *Engine) collectAndBatch( excludedSessionIDs := e.applyIDPrefixToSessionIDs( r.excludedSessionIDs, ) + // A source with no session inside the cwd allow-list must not + // delete archived rows: its exclusions and stale-row cleanup + // would erase sessions whose replacement writes the filter + // vetoes, breaking the ingestion-only contract. Dropping the + // IDs here also keeps them out of parserExcludedIDs, so + // resync's orphan copy still restores the archived rows. + if len(excludedSessionIDs) > 0 && + !e.sourceAllowsParserExclusions(r.processResult) { + excludedSessionIDs = nil + } if len(excludedSessionIDs) > 0 { if _, err := e.db.DeleteParserExcludedSessions( excludedSessionIDs, @@ -3323,6 +3364,21 @@ func (e *Engine) collectAndBatch( } stats.filesOK++ + // Drop sessions outside the cwd allow-list before batching so + // the sync stats can tell an intentionally filtered file apart + // from one whose sessions vanished for an unexplained reason. + // The prepareSessionWrite veto stays as the write-seam backstop. + // Filtered files are deliberately not skip-cached: a later + // allow-list change must be able to pick them up again. + allowed, vetoed := e.splitResultsByCwdFilter(r.results) + stats.cwdFilteredSessions += vetoed + if vetoed > 0 && len(allowed) == 0 { + stats.cwdFilteredFiles++ + progress.SessionsDone++ + e.reportProgress(onProgress, progress) + continue + } + if r.incremental != nil { if err := e.writeIncremental(r.incremental); err != nil { log.Printf("%v", err) @@ -3335,7 +3391,7 @@ func (e *Engine) collectAndBatch( ) stats.messagesIndexed = progress.MessagesIndexed } else { - for _, pr := range r.results { + for _, pr := range allowed { pending = append(pending, pendingWrite{ sess: pr.Session, msgs: pr.Messages, @@ -3359,12 +3415,13 @@ func (e *Engine) collectAndBatch( } if len(pending) >= batchSize { - writtenSessions, writtenMessages, failedWrites := + writtenSessions, writtenMessages, failedWrites, cwdFiltered := e.writeBatch(pending, writeMode, false) stats.RecordSynced(writtenSessions) for range failedWrites { stats.RecordFailed() } + stats.cwdFilteredSessions += cwdFiltered progress.MessagesIndexed += writtenMessages stats.messagesIndexed = progress.MessagesIndexed pending = pending[:0] @@ -3376,12 +3433,13 @@ func (e *Engine) collectAndBatch( flush: if len(pending) > 0 { - writtenSessions, writtenMessages, failedWrites := + writtenSessions, writtenMessages, failedWrites, cwdFiltered := e.writeBatch(pending, writeMode, false) stats.RecordSynced(writtenSessions) for range failedWrites { stats.RecordFailed() } + stats.cwdFilteredSessions += cwdFiltered progress.MessagesIndexed += writtenMessages stats.messagesIndexed = progress.MessagesIndexed } @@ -4875,6 +4933,15 @@ func (e *Engine) tryIncrementalJSONL( return processResult{}, false } + // A session archived before the cwd allow-list was configured + // must not keep growing through the append path, which bypasses + // the prepareSessionWrite veto. Fall back to the full parse path + // so the same veto applies; it also re-derives the cwd from the + // whole file, which covers a stored cwd that predates cwd capture. + if !e.cwdFilter.allows(inc.Cwd) { + return processResult{}, false + } + // Existing rows from an older parser lack new metadata // columns. Force a full parse so the rewrite picks them // up rather than appending new rows on top of stale ones. @@ -5677,17 +5744,20 @@ func (e *Engine) writeBatch( batch []pendingWrite, writeMode syncWriteMode, forceReplace bool, -) (writtenSessions, writtenMessages, failedSessions int) { +) (writtenSessions, writtenMessages, failedSessions, cwdFiltered int) { if writeMode == syncWriteBulk { return e.writeBatchBulk(batch, forceReplace) } resolveWorktreeProject := e.loadWorktreeProjectResolver() for _, pw := range batch { - s, msgs, ok := e.prepareSessionWrite( + s, msgs, verdict := e.prepareSessionWrite( pw, resolveWorktreeProject, ) - if !ok { + if verdict != sessionWriteOK { + if verdict == sessionWriteCwdFiltered { + cwdFiltered++ + } continue } @@ -5794,13 +5864,26 @@ func (e *Engine) writeBatch( writtenSessions++ writtenMessages += len(msgs) } - return writtenSessions, writtenMessages, failedSessions + return writtenSessions, writtenMessages, failedSessions, cwdFiltered } +// sessionWriteVerdict says whether prepareSessionWrite produced a +// writable session and, when it did not, why. The cwd-filter veto is +// distinguished from archive-preserve vetoes so sync stats can count +// filtered sessions: a resync where every discovered session is +// filtered must read as intentional, not as an empty rebuild. +type sessionWriteVerdict int + +const ( + sessionWriteOK sessionWriteVerdict = iota + sessionWritePreserved + sessionWriteCwdFiltered +) + func (e *Engine) prepareSessionWrite( pw pendingWrite, resolveWorktreeProject worktreeProjectResolver, -) (db.Session, []db.Message, bool) { +) (db.Session, []db.Message, sessionWriteVerdict) { msgs := toDBMessages(pw, e.blockedResultCategories) s := toDBSession(pw) applySessionMessageDerivedFields(&s, msgs) @@ -5813,16 +5896,23 @@ func (e *Engine) prepareSessionWrite( } } + // Veto sessions outside the configured cwd allow-list before any + // preserve/merge handling so a filtered session is not written by + // any downstream path. + if !e.cwdFilter.allows(s.Cwd) { + return db.Session{}, nil, sessionWriteCwdFiltered + } + if e.shouldPreserveOpenCodeFormatArchive( pw.sess.Agent, pw.sess.File.Path, s.ID, pw.sess.File.Mtime, derefString(s.FileHash), msgs, ) { - return db.Session{}, nil, false + return db.Session{}, nil, sessionWritePreserved } if mergedMsgs, preserve, archived := e.reconcileVisualStudioCopilotArchive( pw.sess.Agent, s.ID, pw.sess.File.Size, msgs, ); preserve { - return db.Session{}, nil, false + return db.Session{}, nil, sessionWritePreserved } else if mergedMsgs != nil { parsedMsgs := msgs msgs = mergedMsgs @@ -5901,7 +5991,7 @@ func (e *Engine) prepareSessionWrite( _, _, p, h := usageEventTokenTotals(pw.usageEvents, true) s.PeakContextTokens, s.HasPeakContextTokens = p, h } - return s, msgs, true + return s, msgs, sessionWriteOK } func applySessionMessageDerivedFields(s *db.Session, msgs []db.Message) { @@ -6577,18 +6667,21 @@ type projectIdentityCacheEntry struct { func (e *Engine) writeBatchBulk( batch []pendingWrite, forceReplace bool, -) (writtenSessions, writtenMessages, failedSessions int) { +) (writtenSessions, writtenMessages, failedSessions, cwdFiltered int) { writes := make([]db.SessionBatchWrite, 0, len(batch)) sources := make(map[string]batchSourceFile, len(batch)) resolveWorktreeProject := e.loadWorktreeProjectResolver() for _, pw := range batch { tPrep := time.Now() - s, msgs, ok := e.prepareSessionWrite( + s, msgs, verdict := e.prepareSessionWrite( pw, resolveWorktreeProject, ) e.phaseStats.PrepNanos.Add(int64(time.Since(tPrep))) - if !ok { + if verdict != sessionWriteOK { + if verdict == sessionWriteCwdFiltered { + cwdFiltered++ + } continue } replaceMessages := shouldReplaceFullParseMessages( @@ -6618,7 +6711,7 @@ func (e *Engine) writeBatchBulk( } } if len(writes) == 0 { - return 0, 0, 0 + return 0, 0, 0, cwdFiltered } tWrite := time.Now() @@ -6629,7 +6722,7 @@ func (e *Engine) writeBatchBulk( e.phaseStats.BatchedWrites.Add(int64(result.WrittenSessions)) if err != nil { log.Printf("write session batch: %v", err) - return 0, 0, len(writes) + return 0, 0, len(writes), cwdFiltered } for _, id := range result.ExcludedIDs { if source, ok := sources[id]; ok && source.path != "" { @@ -6643,7 +6736,8 @@ func (e *Engine) writeBatchBulk( } return result.WrittenSessions, result.WrittenMessages, - result.FailedSessions + result.FailedSessions, + cwdFiltered } func identityObservationOrZero( @@ -6929,6 +7023,20 @@ func shouldReplaceFullParseMessages( func (e *Engine) writeIncremental( inc *incrementalUpdate, ) error { + // The full path vetoes filtered sessions in prepareSessionWrite; + // this is the equivalent veto at the incremental write seam, so + // no producer can append to a session outside the cwd allow-list. + // tryIncrementalJSONL already refuses such sessions — this guard + // keeps the seam safe for any future producer. + if !e.cwdFilter.allows(inc.cwd) { + log.Printf( + "incremental %s: cwd %q outside the configured "+ + "allow-list, skipping append", + inc.sessionID, inc.cwd, + ) + return nil + } + dbMsgs := toDBMessages( pendingWrite{ sess: parser.ParsedSession{ID: inc.sessionID}, @@ -7100,10 +7208,10 @@ func (e *Engine) writeSessionFullWithResolver( pw pendingWrite, resolveWorktreeProject worktreeProjectResolver, ) error { - s, msgs, ok := e.prepareSessionWrite( + s, msgs, verdict := e.prepareSessionWrite( pw, resolveWorktreeProject, ) - if !ok { + if verdict != sessionWriteOK { return errSessionPreserved } if err := e.db.UpsertSession(s); err != nil { @@ -8285,10 +8393,12 @@ func (e *Engine) SyncSingleSessionContext( // from its directory-name fallback ID to the canonical // meta.json ID and returns the stale fallback ID here; without // this delete a single-session resync would leave both rows in - // the DB and double-count messages and usage. + // the DB and double-count messages and usage. Like + // collectAndBatch, exclusions from a source with no session + // inside the cwd allow-list are frozen so archived rows survive. if excluded := e.applyIDPrefixToSessionIDs( res.excludedSessionIDs, - ); len(excluded) > 0 { + ); len(excluded) > 0 && e.sourceAllowsParserExclusions(res) { if _, err := e.db.DeleteParserExcludedSessions( excluded, ); err != nil { diff --git a/internal/sync/engine_test.go b/internal/sync/engine_test.go index 9d3c49bd1..bdf3c7a1e 100644 --- a/internal/sync/engine_test.go +++ b/internal/sync/engine_test.go @@ -966,7 +966,7 @@ func TestWriteBatchRemoteIDPrefixUsageEvents(t *testing.T) { }}, } - written, _, failed := e.writeBatch( + written, _, failed, _ := e.writeBatch( []pendingWrite{pw}, syncWriteDefault, false, ) require.Equal(t, 0, failed, "no session writes may fail") @@ -996,7 +996,7 @@ func TestProjectIdentityWriteBatchDiscoversLocalGitRemote(t *testing.T) { require.NoError(t, os.Mkdir(cwd, 0o755)) e := NewEngine(database, EngineConfig{Machine: "laptop"}) - written, _, failed := e.writeBatch([]pendingWrite{{ + written, _, failed, _ := e.writeBatch([]pendingWrite{{ sess: parser.ParsedSession{ ID: "identity-local", Project: "repo", @@ -1156,7 +1156,7 @@ func TestProjectIdentityWriteBatchRejectsNonNativeWindowsDriveGitRemote(t *testi require.NoError(t, os.MkdirAll(cwd, 0o755)) e := NewEngine(database, EngineConfig{Machine: "windows-host"}) - written, _, failed := e.writeBatch([]pendingWrite{{ + written, _, failed, _ := e.writeBatch([]pendingWrite{{ sess: parser.ParsedSession{ ID: "identity-windows", Project: "windows", @@ -1191,7 +1191,7 @@ func TestProjectIdentityRemoteWriteSkipsLiveDiscovery(t *testing.T) { return "remote-host:" + path }, }) - written, _, failed := e.writeBatch([]pendingWrite{{ + written, _, failed, _ := e.writeBatch([]pendingWrite{{ sess: parser.ParsedSession{ ID: "identity-remote", Project: "remote-project", @@ -1346,13 +1346,13 @@ func TestWriteBatchAntigravityReplacesMessages(t *testing.T) { } } - written, _, failed := e.writeBatch( + written, _, failed, _ := e.writeBatch( []pendingWrite{mkWrite(false)}, syncWriteDefault, false, ) require.Equal(t, 0, failed) require.Equal(t, 1, written) - written, _, failed = e.writeBatch( + written, _, failed, _ = e.writeBatch( []pendingWrite{mkWrite(true)}, syncWriteDefault, false, ) require.Equal(t, 0, failed) @@ -1402,13 +1402,13 @@ func TestWriteBatchQwenPawReplacesMessages(t *testing.T) { } } - written, _, failed := e.writeBatch( + written, _, failed, _ := e.writeBatch( []pendingWrite{mkWrite("old content")}, syncWriteDefault, false, ) require.Equal(t, 0, failed) require.Equal(t, 1, written) - written, _, failed = e.writeBatch( + written, _, failed, _ = e.writeBatch( []pendingWrite{mkWrite("new content")}, syncWriteDefault, false, ) require.Equal(t, 0, failed) @@ -1533,7 +1533,7 @@ func TestProcessAntigravityWALOnlyUpdateNotSkipped(t *testing.T) { usageEvents: res.results[0].UsageEvents, forceReplace: res.forceReplace, } - written, _, failed := e.writeBatch( + written, _, failed, _ := e.writeBatch( []pendingWrite{pw}, syncWriteDefault, false, ) require.Equal(t, 0, failed) @@ -1662,7 +1662,7 @@ func TestProcessAntigravityBrainOnlyUpdateNotSkipped(t *testing.T) { usageEvents: res.results[0].UsageEvents, forceReplace: res.forceReplace, } - written, _, failed := e.writeBatch( + written, _, failed, _ := e.writeBatch( []pendingWrite{pw}, syncWriteDefault, false, ) require.Equal(t, 0, failed) @@ -2773,10 +2773,10 @@ func TestPrepareSessionWriteReclampsMessageDerivedTokenTotals(t *testing.T) { sess.PeakContextTokens = 999_999_999 sess.HasPeakContextTokens = true - prepared, dbMsgs, ok := e.prepareSessionWrite( + prepared, dbMsgs, verdict := e.prepareSessionWrite( pendingWrite{sess: sess, msgs: msgs}, nil, ) - require.True(t, ok) + require.Equal(t, sessionWriteOK, verdict) require.Len(t, dbMsgs, 3) // The corrupt message row is clamped to the per-message bound. @@ -2804,10 +2804,10 @@ func TestPrepareSessionWriteReclampsMessageDerivedTokenTotals(t *testing.T) { summarySess.PeakContextTokens = summaryPeak summarySess.HasPeakContextTokens = true - preparedSummary, _, ok := e.prepareSessionWrite( + preparedSummary, _, verdict := e.prepareSessionWrite( pendingWrite{sess: summarySess, msgs: msgs}, nil, ) - require.True(t, ok) + require.Equal(t, sessionWriteOK, verdict) assert.Equal(t, summaryTotal, preparedSummary.TotalOutputTokens, "summary-derived total left untouched by per-message clamp") assert.Equal(t, summaryPeak, preparedSummary.PeakContextTokens, @@ -2857,10 +2857,10 @@ func TestPrepareSessionWriteReclampsEventDerivedTokenTotals(t *testing.T) { HasPeakContextTokens: true, } - prepared, _, ok := e.prepareSessionWrite( + prepared, _, verdict := e.prepareSessionWrite( pendingWrite{sess: sess, msgs: msgs, usageEvents: events}, nil, ) - require.True(t, ok) + require.Equal(t, sessionWriteOK, verdict) assert.Equal(t, 1_000_000+1_500_000+maxPlausibleTokens, prepared.TotalOutputTokens, "event-derived total re-derived from clamped usage events") @@ -2904,10 +2904,10 @@ func TestPrepareSessionWritePreservesSummaryUsageEventTokenTotals(t *testing.T) HasPeakContextTokens: true, } - prepared, _, ok := e.prepareSessionWrite( + prepared, _, verdict := e.prepareSessionWrite( pendingWrite{sess: sess, msgs: msgs, usageEvents: events}, nil, ) - require.True(t, ok) + require.Equal(t, sessionWriteOK, verdict) assert.Equal(t, rawTotal, prepared.TotalOutputTokens, "session-summary usage event must not make the session aggregate event-derived") assert.Equal(t, rawPeak, prepared.PeakContextTokens, @@ -2959,10 +2959,10 @@ func TestPrepareSessionWriteReclampsEventDerivedCacheContext(t *testing.T) { HasPeakContextTokens: true, } - prepared, _, ok := e.prepareSessionWrite( + prepared, _, verdict := e.prepareSessionWrite( pendingWrite{sess: sess, msgs: msgs, usageEvents: events}, nil, ) - require.True(t, ok) + require.Equal(t, sessionWriteOK, verdict) assert.Equal(t, 100_000+100_000+maxPlausibleTokens, prepared.TotalOutputTokens, "event-derived total re-derived from clamped output tokens") @@ -3017,10 +3017,10 @@ func TestPrepareSessionWriteReclampsEventDerivedMixedSignTokens(t *testing.T) { HasPeakContextTokens: true, } - prepared, _, ok := e.prepareSessionWrite( + prepared, _, verdict := e.prepareSessionWrite( pendingWrite{sess: sess, msgs: msgs, usageEvents: events}, nil, ) - require.True(t, ok) + require.Equal(t, sessionWriteOK, verdict) // Negative output floors to 0 (dropped); over-bound output clamps to 2M. assert.Equal(t, 1_000_000+maxPlausibleTokens, prepared.TotalOutputTokens, "negative event excluded, over-bound event clamped in event total") @@ -4675,7 +4675,7 @@ func TestWriteIncrementalBlanksImplausibleEndedAt(t *testing.T) { Timestamp: start, }}, } - _, _, failed := e.writeBatch( + _, _, failed, _ := e.writeBatch( []pendingWrite{pw}, syncWriteDefault, false, ) require.Equal(t, 0, failed, "initial session write must not fail") @@ -4747,7 +4747,7 @@ func TestWriteIncrementalKeepsPlausibleEndedAt(t *testing.T) { Timestamp: start, }}, } - _, _, failed := e.writeBatch( + _, _, failed, _ := e.writeBatch( []pendingWrite{pw}, syncWriteDefault, false, ) require.Equal(t, 0, failed, "initial session write must not fail") diff --git a/internal/sync/hermes_archive_test.go b/internal/sync/hermes_archive_test.go index 58026330c..c9e2e5151 100644 --- a/internal/sync/hermes_archive_test.go +++ b/internal/sync/hermes_archive_test.go @@ -203,7 +203,7 @@ func TestProcessFileHermesArchivePersistsAggregateFingerprint(t *testing.T) { usageEvents: result.UsageEvents, }) } - written, _, failed := engine.writeBatch(pending, syncWriteDefault, true) + written, _, failed, _ := engine.writeBatch(pending, syncWriteDefault, true) require.Equal(t, 0, failed) require.NotZero(t, written) diff --git a/internal/sync/parsediff.go b/internal/sync/parsediff.go index 8d4ec7148..c119a7b5d 100644 --- a/internal/sync/parsediff.go +++ b/internal/sync/parsediff.go @@ -660,9 +660,9 @@ func (e *Engine) parseDiffCollectFile( usageEvents: pr.UsageEvents, needsRetry: job.needsRetryForSession(pr.Session.ID), } - prepared, msgs, ok := e.prepareSessionWrite(pw, resolver) + prepared, msgs, verdict := e.prepareSessionWrite(pw, resolver) id := prepared.ID - if !ok { + if verdict != sessionWriteOK { // prepareSessionWrite returns a zero session on veto; // reconstruct the final ID the way applyRemoteRewrites // would have. @@ -674,7 +674,7 @@ func (e *Engine) parseDiffCollectFile( } var fields []FieldDiff - compare := ok && !pw.needsRetry && + compare := verdict == sessionWriteOK && !pw.needsRetry && stored != nil && stored.DeletedAt == nil if compare { events, _ := toDBUsageEvents(id, pw.usageEvents) @@ -770,7 +770,7 @@ func (e *Engine) parseDiffCollectFile( class, reason := classifyParseDiffSession( pw.needsRetry, - ok, + verdict == sessionWriteOK, stored != nil, stored != nil && stored.DeletedAt != nil, stored != nil && stored.DataVersion < db.CurrentDataVersion(), diff --git a/internal/sync/parsediff_compare_test.go b/internal/sync/parsediff_compare_test.go index 3ecc2707f..750e21e2f 100644 --- a/internal/sync/parsediff_compare_test.go +++ b/internal/sync/parsediff_compare_test.go @@ -1061,14 +1061,14 @@ func TestFingerprintTwinMatchesDB(t *testing.T) { }, } - written, _, failed := e.writeBatch( + written, _, failed, _ := e.writeBatch( []pendingWrite{pw}, syncWriteBulk, false, ) require.Equal(t, 1, written, "session must be written") require.Zero(t, failed) - prepared, msgs, ok := e.prepareSessionWrite(pw, nil) - require.True(t, ok) + prepared, msgs, verdict := e.prepareSessionWrite(pw, nil) + require.Equal(t, sessionWriteOK, verdict) require.NotEmpty(t, msgs) storedFP, err := d.MessageTokenFingerprint(prepared.ID) @@ -1163,14 +1163,14 @@ func TestCompareStoredSessionRoundTrip(t *testing.T) { }, } - written, _, failed := e.writeBatch( + written, _, failed, _ := e.writeBatch( []pendingWrite{pw}, syncWriteBulk, false, ) require.Equal(t, 1, written) require.Zero(t, failed) - prepared, msgs, ok := e.prepareSessionWrite(pw, nil) - require.True(t, ok) + prepared, msgs, verdict := e.prepareSessionWrite(pw, nil) + require.Equal(t, sessionWriteOK, verdict) events, _ := toDBUsageEvents(prepared.ID, pw.usageEvents) stored := pdFetchStored(t, d, prepared.ID) @@ -1215,7 +1215,7 @@ func TestCompareStoredSessionDetectsDrift(t *testing.T) { }, }, } - written, _, failed := e.writeBatch( + written, _, failed, _ := e.writeBatch( []pendingWrite{pw}, syncWriteBulk, false, ) require.Equal(t, 1, written) @@ -1223,8 +1223,8 @@ func TestCompareStoredSessionDetectsDrift(t *testing.T) { // Simulate parser drift: the new parse reports a different model. pw.msgs[0].Model = "claude-haiku" - prepared, msgs, ok := e.prepareSessionWrite(pw, nil) - require.True(t, ok) + prepared, msgs, verdict := e.prepareSessionWrite(pw, nil) + require.Equal(t, sessionWriteOK, verdict) stored := pdFetchStored(t, d, prepared.ID) diffs, err := e.compareStoredSession( @@ -1258,7 +1258,7 @@ func pdWriteSingleMessageSession( }, msgs: []parser.ParsedMessage{msg}, } - written, _, failed := e.writeBatch( + written, _, failed, _ := e.writeBatch( []pendingWrite{pw}, syncWriteBulk, false, ) require.Equal(t, 1, written) @@ -1280,8 +1280,8 @@ func TestCompareStoredSessionDetectsContentDrift(t *testing.T) { // New parse: same model and tokens, longer body. pw.msgs[0].Content = "a much longer reply body" pw.msgs[0].ContentLength = len(pw.msgs[0].Content) - prepared, msgs, ok := e.prepareSessionWrite(pw, nil) - require.True(t, ok) + prepared, msgs, verdict := e.prepareSessionWrite(pw, nil) + require.Equal(t, sessionWriteOK, verdict) stored := pdFetchStored(t, d, prepared.ID) diffs, err := e.compareStoredSession( @@ -1306,8 +1306,8 @@ func TestCompareStoredSessionDetectsMetadataDrift(t *testing.T) { // New parse flips only is_sidechain: same model, tokens, content. pw.msgs[0].IsSidechain = true - prepared, msgs, ok := e.prepareSessionWrite(pw, nil) - require.True(t, ok) + prepared, msgs, verdict := e.prepareSessionWrite(pw, nil) + require.Equal(t, sessionWriteOK, verdict) stored := pdFetchStored(t, d, prepared.ID) diffs, err := e.compareStoredSession( @@ -1380,7 +1380,7 @@ func pdWriteToolSession( d := openTestDB(t) e := NewEngine(d, EngineConfig{Machine: "test-machine"}) pw := pdToolSession(id) - written, _, failed := e.writeBatch( + written, _, failed, _ := e.writeBatch( []pendingWrite{pw}, syncWriteBulk, false, ) require.Equal(t, 1, written) @@ -1393,8 +1393,8 @@ func pdWriteToolSession( // way TestFingerprintTwinMatchesDB does for the message fingerprints. func TestToolCallAndFlagsFingerprintTwinsMatchDB(t *testing.T) { e, d, pw := pdWriteToolSession(t, "pd-tool-twin") - prepared, msgs, ok := e.prepareSessionWrite(pw, nil) - require.True(t, ok) + prepared, msgs, verdict := e.prepareSessionWrite(pw, nil) + require.Equal(t, sessionWriteOK, verdict) storedFlagsFP, err := d.MessageFlagsFingerprint(prepared.ID) require.NoError(t, err) @@ -1427,8 +1427,8 @@ func TestToolCallDiffDetectsFilePath(t *testing.T) { // system message must compare identical against itself. func TestCompareStoredSessionRoundTripToolCalls(t *testing.T) { e, d, pw := pdWriteToolSession(t, "pd-tool-rt") - prepared, msgs, ok := e.prepareSessionWrite(pw, nil) - require.True(t, ok) + prepared, msgs, verdict := e.prepareSessionWrite(pw, nil) + require.Equal(t, sessionWriteOK, verdict) stored := pdFetchStored(t, d, prepared.ID) diffs, err := e.compareStoredSession( @@ -1446,8 +1446,8 @@ func TestCompareStoredSessionRoundTripToolCalls(t *testing.T) { func TestCompareStoredSessionDetectsToolCallDrift(t *testing.T) { e, d, pw := pdWriteToolSession(t, "pd-tool-drift") pw.msgs[1].ToolCalls[0].ToolName = "Grep" - prepared, msgs, ok := e.prepareSessionWrite(pw, nil) - require.True(t, ok) + prepared, msgs, verdict := e.prepareSessionWrite(pw, nil) + require.Equal(t, sessionWriteOK, verdict) stored := pdFetchStored(t, d, prepared.ID) diffs, err := e.compareStoredSession( @@ -1465,8 +1465,8 @@ func TestCompareStoredSessionDetectsToolCallDrift(t *testing.T) { func TestCompareStoredSessionDetectsFlagDrift(t *testing.T) { e, d, pw := pdWriteToolSession(t, "pd-flag-drift") pw.msgs[1].HasThinking = false - prepared, msgs, ok := e.prepareSessionWrite(pw, nil) - require.True(t, ok) + prepared, msgs, verdict := e.prepareSessionWrite(pw, nil) + require.Equal(t, sessionWriteOK, verdict) stored := pdFetchStored(t, d, prepared.ID) diffs, err := e.compareStoredSession( diff --git a/internal/sync/progress.go b/internal/sync/progress.go index c69b3c105..cbaef0df6 100644 --- a/internal/sync/progress.go +++ b/internal/sync/progress.go @@ -76,6 +76,16 @@ type SyncStats struct { messagesIndexed int // unexported: progress message counter parserExcludedFiles int // file-level intentional parser exclusions parserExcludedIDs []string + // cwdFilteredSessions counts sessions vetoed by the + // sync_include_cwd_prefixes allow-list. The resync abort guard uses + // it so a run where every discovered session is filtered reads as + // an intentional result rather than an unsafe empty rebuild. + cwdFilteredSessions int + // cwdFilteredFiles counts files whose every parsed session was + // vetoed by the allow-list. Together with parserExcludedFiles it + // must account for all of filesOK before the resync abort guard + // treats a zero-write run as intentional. + cwdFilteredFiles int } // AnomalyStats aggregates parser-output anomaly signals observed during a diff --git a/internal/sync/provider_process_test.go b/internal/sync/provider_process_test.go index 064e9ba84..9ddfbff24 100644 --- a/internal/sync/provider_process_test.go +++ b/internal/sync/provider_process_test.go @@ -141,7 +141,7 @@ func TestProcessFileProviderSkipsStoredFreshSource(t *testing.T) { }) require.NoError(t, first.err) require.Len(t, first.results, 1) - written, _, failed := engine.writeBatch( + written, _, failed, _ := engine.writeBatch( []pendingWrite{{ sess: first.results[0].Session, msgs: first.results[0].Messages, @@ -224,7 +224,7 @@ func TestProcessFileProviderPiebaldSkipsStoredFreshSource(t *testing.T) { }) require.NoError(t, first.err) require.Len(t, first.results, 1) - written, _, failed := engine.writeBatch( + written, _, failed, _ := engine.writeBatch( []pendingWrite{{ sess: first.results[0].Session, msgs: first.results[0].Messages, @@ -617,7 +617,7 @@ func TestProcessFileProviderDevinSkipsStoredFreshSource(t *testing.T) { require.NotZero(t, storedMtime) assert.GreaterOrEqual(t, storedMtime, transcriptProcessProviderMtime(t, transcriptPath)) - written, _, failed := engine.writeBatch( + written, _, failed, _ := engine.writeBatch( []pendingWrite{{ sess: first.results[0].Session, msgs: first.results[0].Messages, @@ -672,7 +672,7 @@ func TestProcessFileProviderDevinReparsesTranscriptOnlyChange(t *testing.T) { }) require.NoError(t, first.err) require.Len(t, first.results, 1) - written, _, failed := engine.writeBatch( + written, _, failed, _ := engine.writeBatch( []pendingWrite{{ sess: first.results[0].Session, msgs: first.results[0].Messages, @@ -747,7 +747,7 @@ func TestProcessFileProviderDevinSameSizeSameMtimeTranscriptRewriteReparses(t *t require.NotZero(t, initialMtime) require.NotEmpty(t, initialHash) - written, _, failed := engine.writeBatch( + written, _, failed, _ := engine.writeBatch( []pendingWrite{{ sess: first.results[0].Session, msgs: first.results[0].Messages, @@ -947,7 +947,7 @@ func writeProcessProviderDevinResult( ) { t.Helper() require.Len(t, result.results, 1) - written, _, failed := engine.writeBatch( + written, _, failed, _ := engine.writeBatch( []pendingWrite{{ sess: result.results[0].Session, msgs: result.results[0].Messages, diff --git a/internal/sync/qoder_test.go b/internal/sync/qoder_test.go index a5b7450d4..9062163f1 100644 --- a/internal/sync/qoder_test.go +++ b/internal/sync/qoder_test.go @@ -239,7 +239,7 @@ func writeProcessQoderResult( result processResult, ) { t.Helper() - written, _, failed := engine.writeBatch( + written, _, failed, _ := engine.writeBatch( []pendingWrite{{ sess: result.results[0].Session, msgs: result.results[0].Messages, diff --git a/internal/sync/s3_provider_discovery_test.go b/internal/sync/s3_provider_discovery_test.go index cf786cec9..1e0533143 100644 --- a/internal/sync/s3_provider_discovery_test.go +++ b/internal/sync/s3_provider_discovery_test.go @@ -72,7 +72,7 @@ func TestProcessFileS3ProviderDiscoveredRoutesToS3Path(t *testing.T) { require.NoError(t, res.err) require.Len(t, res.results, 1) - written, _, failed := e.writeBatch([]pendingWrite{{ + written, _, failed, _ := e.writeBatch([]pendingWrite{{ sess: res.results[0].Session, msgs: res.results[0].Messages, }}, syncWriteDefault, false) diff --git a/internal/sync/s3_sidecars_test.go b/internal/sync/s3_sidecars_test.go index e0210d6ab..018b26783 100644 --- a/internal/sync/s3_sidecars_test.go +++ b/internal/sync/s3_sidecars_test.go @@ -227,7 +227,7 @@ func TestProcessS3ClaudeHydratedSidecarReplacesStoredPreview(t *testing.T) { SourceMtime: time.Date(2026, 6, 24, 12, 14, 0, 0, time.UTC).UnixNano(), }) require.NoError(t, first.err) - written, _, failed := e.writeBatch([]pendingWrite{{ + written, _, failed, _ := e.writeBatch([]pendingWrite{{ sess: first.results[0].Session, msgs: first.results[0].Messages, }}, syncWriteDefault, false) @@ -245,7 +245,7 @@ func TestProcessS3ClaudeHydratedSidecarReplacesStoredPreview(t *testing.T) { }) require.NoError(t, second.err) require.True(t, second.forceReplace) - written, _, failed = e.writeBatch([]pendingWrite{{ + written, _, failed, _ = e.writeBatch([]pendingWrite{{ sess: second.results[0].Session, msgs: second.results[0].Messages, forceReplace: second.forceReplace, @@ -346,7 +346,7 @@ func TestSyncSingleSessionS3ClaudeSidecarOnlyChangeReplacesPreview( SourceMtime: transcriptMtime.UnixNano(), }) require.NoError(t, first.err) - written, _, failed := e.writeBatch([]pendingWrite{{ + written, _, failed, _ := e.writeBatch([]pendingWrite{{ sess: first.results[0].Session, msgs: first.results[0].Messages, }}, syncWriteDefault, false) @@ -412,7 +412,7 @@ func TestProcessS3ClaudeMissingSidecarReplacesStoredHydratedOutput( SourceMtime: time.Date(2026, 6, 24, 12, 17, 0, 0, time.UTC).UnixNano(), }) require.NoError(t, first.err) - written, _, failed := e.writeBatch([]pendingWrite{{ + written, _, failed, _ := e.writeBatch([]pendingWrite{{ sess: first.results[0].Session, msgs: first.results[0].Messages, forceReplace: first.forceReplace, @@ -431,7 +431,7 @@ func TestProcessS3ClaudeMissingSidecarReplacesStoredHydratedOutput( }) require.NoError(t, second.err) require.True(t, second.forceReplace) - written, _, failed = e.writeBatch([]pendingWrite{{ + written, _, failed, _ = e.writeBatch([]pendingWrite{{ sess: second.results[0].Session, msgs: second.results[0].Messages, forceReplace: second.forceReplace, diff --git a/internal/sync/s3_source_test.go b/internal/sync/s3_source_test.go index 39e271e87..5f1d32731 100644 --- a/internal/sync/s3_source_test.go +++ b/internal/sync/s3_source_test.go @@ -154,7 +154,7 @@ func TestProcessFileS3ChangedFingerprintReplacesStoredMessages(t *testing.T) { }) require.NoError(t, res.err) require.Len(t, res.results, 1) - written, _, failed := e.writeBatch([]pendingWrite{{ + written, _, failed, _ := e.writeBatch([]pendingWrite{{ sess: res.results[0].Session, msgs: res.results[0].Messages, forceReplace: res.forceReplace, @@ -175,7 +175,7 @@ func TestProcessFileS3ChangedFingerprintReplacesStoredMessages(t *testing.T) { require.NoError(t, res.err) require.False(t, res.skip) require.Len(t, res.results, 1) - written, _, failed = e.writeBatch([]pendingWrite{{ + written, _, failed, _ = e.writeBatch([]pendingWrite{{ sess: res.results[0].Session, msgs: res.results[0].Messages, forceReplace: res.forceReplace, @@ -847,7 +847,7 @@ func TestProcessFileS3SameMetadataDifferentURIRewritesSourcePath(t *testing.T) { require.True(t, fetched) require.Len(t, res.results, 1) - written, _, failed := e.writeBatch([]pendingWrite{{ + written, _, failed, _ := e.writeBatch([]pendingWrite{{ sess: res.results[0].Session, msgs: res.results[0].Messages, }}, syncWriteDefault, false) @@ -1002,7 +1002,7 @@ func TestSyncSingleSessionS3PreservesStoredMachine(t *testing.T) { }) require.NoError(t, res.err) require.Len(t, res.results, 1) - written, _, failed := e.writeBatch([]pendingWrite{{ + written, _, failed, _ := e.writeBatch([]pendingWrite{{ sess: res.results[0].Session, msgs: res.results[0].Messages, }}, syncWriteDefault, false) @@ -1083,7 +1083,7 @@ func TestSyncSingleSessionS3WithoutMachineNamespaceUpdatesRawID( }) require.NoError(t, res.err) require.Len(t, res.results, 1) - written, _, failed := e.writeBatch([]pendingWrite{{ + written, _, failed, _ := e.writeBatch([]pendingWrite{{ sess: res.results[0].Session, msgs: res.results[0].Messages, }}, syncWriteDefault, false) diff --git a/internal/sync/s3_test.go b/internal/sync/s3_test.go index bc54f05e1..2b270682c 100644 --- a/internal/sync/s3_test.go +++ b/internal/sync/s3_test.go @@ -43,7 +43,7 @@ func TestProcessS3SessionNamespacesIDsBySourceMachine(t *testing.T) { require.NoError(t, res.err) require.Len(t, res.results, 1) - written, _, failed := e.writeBatch([]pendingWrite{{ + written, _, failed, _ := e.writeBatch([]pendingWrite{{ sess: res.results[0].Session, msgs: res.results[0].Messages, }}, syncWriteDefault, false) @@ -89,7 +89,7 @@ func TestProcessS3CodexNamespacesIDsBySourceMachine(t *testing.T) { require.NoError(t, res.err) require.Len(t, res.results, 1) - written, _, failed := e.writeBatch([]pendingWrite{{ + written, _, failed, _ := e.writeBatch([]pendingWrite{{ sess: res.results[0].Session, msgs: res.results[0].Messages, }}, syncWriteDefault, false) @@ -490,7 +490,7 @@ func TestProcessS3ClaudeSubagentPreservesParentLayout(t *testing.T) { require.NoError(t, res.err) require.Len(t, res.results, 1) - written, _, failed := e.writeBatch([]pendingWrite{{ + written, _, failed, _ := e.writeBatch([]pendingWrite{{ sess: res.results[0].Session, msgs: res.results[0].Messages, }}, syncWriteDefault, false)