From 48d829e1dc76462726a71451c09a65744fc0492a Mon Sep 17 00:00:00 2001 From: akarineren Date: Mon, 27 Jul 2026 11:56:35 +0900 Subject: [PATCH 1/5] Sort usage entries by timestamp in descending order across providers and update related comments for clarity --- internal/agentusage/providers.go | 113 ++++++++++++++++++++++++------- internal/claudeusage/loader.go | 1 + internal/claudeusage/provider.go | 10 ++- internal/codexusage/provider.go | 13 +++- internal/usagedb/db.go | 4 +- internal/usagedb/db_test.go | 4 +- 6 files changed, 111 insertions(+), 34 deletions(-) diff --git a/internal/agentusage/providers.go b/internal/agentusage/providers.go index 4c2c52f..a290431 100644 --- a/internal/agentusage/providers.go +++ b/internal/agentusage/providers.go @@ -1,10 +1,19 @@ package agentusage import ( + "sort" + "github.com/tokitoki-dev/tokitoki-cli/internal/usage" "github.com/tokitoki-dev/tokitoki-cli/internal/usageprovider" ) +func sortEntriesByTimestampDesc(entries []usage.Entry) []usage.Entry { + sort.Slice(entries, func(i, j int) bool { + return entries[i].Timestamp.After(entries[j].Timestamp) + }) + return entries +} + // providerBase carries the scan configuration shared by every agent // provider. filter skips source files whose events are already ingested; // the SQLite-backed providers (Kilo, Hermes, Goose) deliberately do not @@ -89,9 +98,13 @@ func (p CopilotProvider) WithFileFilter(filter usage.FileFilter) usageprovider.P return p } -// Entries loads normalized GitHub Copilot CLI usage entries. +// Entries loads normalized GitHub Copilot CLI usage entries, newest first. func (p CopilotProvider) Entries() ([]usage.Entry, error) { - return loadCopilotEntries(p.paths, p.filter) + entries, err := loadCopilotEntries(p.paths, p.filter) + if err != nil { + return nil, err + } + return sortEntriesByTimestampDesc(entries), nil } // WithPaths returns a Gemini CLI provider configured with data roots. @@ -109,9 +122,13 @@ func (p GeminiProvider) WithFileFilter(filter usage.FileFilter) usageprovider.Pr return p } -// Entries loads normalized Gemini CLI usage entries. +// Entries loads normalized Gemini CLI usage entries, newest first. func (p GeminiProvider) Entries() ([]usage.Entry, error) { - return loadGeminiEntries(p.paths, p.filter) + entries, err := loadGeminiEntries(p.paths, p.filter) + if err != nil { + return nil, err + } + return sortEntriesByTimestampDesc(entries), nil } // WithPaths returns a Kimi provider configured with data roots. @@ -129,9 +146,13 @@ func (p KimiProvider) WithFileFilter(filter usage.FileFilter) usageprovider.Prov return p } -// Entries loads normalized Kimi usage entries. +// Entries loads normalized Kimi usage entries, newest first. func (p KimiProvider) Entries() ([]usage.Entry, error) { - return loadKimiEntries(p.paths, p.filter) + entries, err := loadKimiEntries(p.paths, p.filter) + if err != nil { + return nil, err + } + return sortEntriesByTimestampDesc(entries), nil } // WithPaths returns a Qwen provider configured with data roots. @@ -149,9 +170,13 @@ func (p QwenProvider) WithFileFilter(filter usage.FileFilter) usageprovider.Prov return p } -// Entries loads normalized Qwen usage entries. +// Entries loads normalized Qwen usage entries, newest first. func (p QwenProvider) Entries() ([]usage.Entry, error) { - return loadQwenEntries(p.paths, p.filter) + entries, err := loadQwenEntries(p.paths, p.filter) + if err != nil { + return nil, err + } + return sortEntriesByTimestampDesc(entries), nil } // WithPaths returns an OpenClaw provider configured with data roots. @@ -169,9 +194,13 @@ func (p OpenClawProvider) WithFileFilter(filter usage.FileFilter) usageprovider. return p } -// Entries loads normalized OpenClaw usage entries. +// Entries loads normalized OpenClaw usage entries, newest first. func (p OpenClawProvider) Entries() ([]usage.Entry, error) { - return loadOpenClawEntries(p.paths, p.filter) + entries, err := loadOpenClawEntries(p.paths, p.filter) + if err != nil { + return nil, err + } + return sortEntriesByTimestampDesc(entries), nil } // WithPaths returns a pi-agent provider configured with data roots. @@ -189,9 +218,13 @@ func (p PiProvider) WithFileFilter(filter usage.FileFilter) usageprovider.Provid return p } -// Entries loads normalized pi-agent usage entries. +// Entries loads normalized pi-agent usage entries, newest first. func (p PiProvider) Entries() ([]usage.Entry, error) { - return loadPiEntries(p.paths, p.filter) + entries, err := loadPiEntries(p.paths, p.filter) + if err != nil { + return nil, err + } + return sortEntriesByTimestampDesc(entries), nil } // WithPaths returns an Amp provider configured with data roots. @@ -209,9 +242,13 @@ func (p AmpProvider) WithFileFilter(filter usage.FileFilter) usageprovider.Provi return p } -// Entries loads normalized Amp usage entries. +// Entries loads normalized Amp usage entries, newest first. func (p AmpProvider) Entries() ([]usage.Entry, error) { - return loadAmpEntries(p.paths, p.filter) + entries, err := loadAmpEntries(p.paths, p.filter) + if err != nil { + return nil, err + } + return sortEntriesByTimestampDesc(entries), nil } // WithPaths returns a Droid provider configured with data roots. @@ -229,9 +266,13 @@ func (p DroidProvider) WithFileFilter(filter usage.FileFilter) usageprovider.Pro return p } -// Entries loads normalized Droid usage entries. +// Entries loads normalized Droid usage entries, newest first. func (p DroidProvider) Entries() ([]usage.Entry, error) { - return loadDroidEntries(p.paths, p.filter) + entries, err := loadDroidEntries(p.paths, p.filter) + if err != nil { + return nil, err + } + return sortEntriesByTimestampDesc(entries), nil } // WithPaths returns a Kilo provider configured with data roots. @@ -242,9 +283,13 @@ func (KiloProvider) WithPaths(paths []string) usageprovider.Provider { // Provider returns the Kilo provider id. func (KiloProvider) Provider() usage.Provider { return usage.ProviderKilo } -// Entries loads normalized Kilo usage entries. +// Entries loads normalized Kilo usage entries, newest first. func (p KiloProvider) Entries() ([]usage.Entry, error) { - return loadKiloEntries(p.paths) + entries, err := loadKiloEntries(p.paths) + if err != nil { + return nil, err + } + return sortEntriesByTimestampDesc(entries), nil } // WithPaths returns a Hermes provider configured with data roots. @@ -255,9 +300,13 @@ func (HermesProvider) WithPaths(paths []string) usageprovider.Provider { // Provider returns the Hermes provider id. func (HermesProvider) Provider() usage.Provider { return usage.ProviderHermes } -// Entries loads normalized Hermes Agent usage entries. +// Entries loads normalized Hermes Agent usage entries, newest first. func (p HermesProvider) Entries() ([]usage.Entry, error) { - return loadHermesEntries(p.paths) + entries, err := loadHermesEntries(p.paths) + if err != nil { + return nil, err + } + return sortEntriesByTimestampDesc(entries), nil } // WithPaths returns a Codebuff provider configured with data roots. @@ -275,9 +324,13 @@ func (p CodebuffProvider) WithFileFilter(filter usage.FileFilter) usageprovider. return p } -// Entries loads normalized Codebuff usage entries. +// Entries loads normalized Codebuff usage entries, newest first. func (p CodebuffProvider) Entries() ([]usage.Entry, error) { - return loadCodebuffEntries(p.paths, p.filter) + entries, err := loadCodebuffEntries(p.paths, p.filter) + if err != nil { + return nil, err + } + return sortEntriesByTimestampDesc(entries), nil } // WithPaths returns an OpenCode provider configured with data roots. @@ -296,9 +349,13 @@ func (p OpenCodeProvider) WithFileFilter(filter usage.FileFilter) usageprovider. return p } -// Entries loads normalized OpenCode usage entries. +// Entries loads normalized OpenCode usage entries, newest first. func (p OpenCodeProvider) Entries() ([]usage.Entry, error) { - return loadOpenCodeEntries(p.paths, p.filter) + entries, err := loadOpenCodeEntries(p.paths, p.filter) + if err != nil { + return nil, err + } + return sortEntriesByTimestampDesc(entries), nil } // WithPaths returns a Goose provider configured with data roots. @@ -309,7 +366,11 @@ func (GooseProvider) WithPaths(paths []string) usageprovider.Provider { // Provider returns the Goose provider id. func (GooseProvider) Provider() usage.Provider { return usage.ProviderGoose } -// Entries loads normalized Goose usage entries. +// Entries loads normalized Goose usage entries, newest first. func (p GooseProvider) Entries() ([]usage.Entry, error) { - return loadGooseEntries(p.paths) + entries, err := loadGooseEntries(p.paths) + if err != nil { + return nil, err + } + return sortEntriesByTimestampDesc(entries), nil } diff --git a/internal/claudeusage/loader.go b/internal/claudeusage/loader.go index 00eae27..85520bb 100644 --- a/internal/claudeusage/loader.go +++ b/internal/claudeusage/loader.go @@ -756,3 +756,4 @@ func pathParts(path string) []string { }) return parts } + diff --git a/internal/claudeusage/provider.go b/internal/claudeusage/provider.go index 27fa81c..6a96afc 100644 --- a/internal/claudeusage/provider.go +++ b/internal/claudeusage/provider.go @@ -1,6 +1,8 @@ package claudeusage import ( + "sort" + "github.com/tokitoki-dev/tokitoki-cli/internal/usage" "github.com/tokitoki-dev/tokitoki-cli/internal/usageprovider" ) @@ -31,11 +33,15 @@ func (Provider) Provider() usage.Provider { return usage.ProviderClaude } -// Entries loads normalized Claude usage entries. +// Entries loads normalized Claude usage entries, newest first. func (p Provider) Entries() ([]usage.Entry, error) { entries, err := LoadEntriesFromPaths(p.paths, "", p.filter) if err != nil { return nil, err } - return ConvertEntries(entries), nil + converted := ConvertEntries(entries) + sort.Slice(converted, func(i, j int) bool { + return converted[i].Timestamp.After(converted[j].Timestamp) + }) + return converted, nil } diff --git a/internal/codexusage/provider.go b/internal/codexusage/provider.go index fad493b..02debf7 100644 --- a/internal/codexusage/provider.go +++ b/internal/codexusage/provider.go @@ -1,6 +1,8 @@ package codexusage import ( + "sort" + "github.com/tokitoki-dev/tokitoki-cli/internal/usage" "github.com/tokitoki-dev/tokitoki-cli/internal/usageprovider" ) @@ -31,7 +33,14 @@ func (Provider) Provider() usage.Provider { return usage.ProviderCodex } -// Entries loads normalized Codex usage entries. +// Entries loads normalized Codex usage entries, newest first. func (p Provider) Entries() ([]usage.Entry, error) { - return LoadEntriesFromPaths(p.paths, "", p.filter) + entries, err := LoadEntriesFromPaths(p.paths, "", p.filter) + if err != nil { + return nil, err + } + sort.Slice(entries, func(i, j int) bool { + return entries[i].Timestamp.After(entries[j].Timestamp) + }) + return entries, nil } diff --git a/internal/usagedb/db.go b/internal/usagedb/db.go index 94d6436..c541976 100644 --- a/internal/usagedb/db.go +++ b/internal/usagedb/db.go @@ -160,7 +160,7 @@ func (s *DB) UpsertScannedFiles(states map[string]FileState) error { return tx.Commit() } -// PendingEvents returns events due for upload at now, oldest first. A limit +// PendingEvents returns events due for upload at now, newest first. A limit // of zero or less means no limit. func (s *DB) PendingEvents(now time.Time, limit int) ([]usage.Entry, error) { if limit <= 0 { @@ -169,7 +169,7 @@ func (s *DB) PendingEvents(now time.Time, limit int) ([]usage.Entry, error) { rows, err := s.db.Query(` SELECT payload FROM usage_events WHERE status IN ('pending', 'failed') AND next_attempt_at <= ? - ORDER BY ts, id + ORDER BY ts DESC, id LIMIT ?`, now.Unix(), limit) if err != nil { return nil, err diff --git a/internal/usagedb/db_test.go b/internal/usagedb/db_test.go index 1b98b4a..79639e4 100644 --- a/internal/usagedb/db_test.go +++ b/internal/usagedb/db_test.go @@ -152,8 +152,8 @@ func TestPendingEventsOrdersByTimestampAndHonorsLimit(t *testing.T) { if err != nil { t.Fatal(err) } - if len(pending) != 1 || pending[0].ID != "event-old" { - t.Fatalf("pending = %+v, want oldest event first", pending) + if len(pending) != 1 || pending[0].ID != "event-new" { + t.Fatalf("pending = %+v, want newest event first", pending) } } From 7fe3605fd1d8d44eb1968977c929eee37437cd22 Mon Sep 17 00:00:00 2001 From: akarineren Date: Mon, 27 Jul 2026 12:24:02 +0900 Subject: [PATCH 2/5] Refactor data storage structure to use hierarchical directories and migrate existing files --- internal/store/lock.go | 8 ++- internal/store/store.go | 97 +++++++++++++++++++++++++++++++++--- internal/store/store_test.go | 14 ++++-- internal/usagedb/db.go | 5 ++ pkg/agentlib/agentlib.go | 4 +- 5 files changed, 113 insertions(+), 15 deletions(-) diff --git a/internal/store/lock.go b/internal/store/lock.go index a6e71a6..de91487 100644 --- a/internal/store/lock.go +++ b/internal/store/lock.go @@ -32,10 +32,14 @@ func AcquireDataLock(dir string, timeout time.Duration) (*DataLock, error) { return lock, nil } -// AcquireLock takes an exclusive advisory lock on dir/name, waiting up to +// AcquireLock takes an exclusive advisory lock on dir/state/name, waiting up to // timeout before giving up with ErrLockBusy. func AcquireLock(dir, name string, timeout time.Duration) (*DataLock, error) { - path := filepath.Join(dir, name) + stateDir := filepath.Join(dir, "state") + if err := os.MkdirAll(stateDir, 0o700); err != nil { + return nil, fmt.Errorf("create state dir: %w", err) + } + path := filepath.Join(stateDir, name) file, err := os.OpenFile(path, os.O_CREATE|os.O_RDWR, 0o600) if err != nil { return nil, fmt.Errorf("open lock file: %w", err) diff --git a/internal/store/store.go b/internal/store/store.go index 554cff2..c2f0ad4 100644 --- a/internal/store/store.go +++ b/internal/store/store.go @@ -12,8 +12,12 @@ import ( ) const ( - UsageDBFile = "usage.db" - apiKeyFile = "api_key" + configDirName = "config" + dataDirName = "data" + stateDirName = "state" + + UsageDBFile = "tokitoki.db" + apiKeyFile = "api_key" directoryMod = 0o700 apiKeyFileMod = 0o600 ) @@ -39,6 +43,12 @@ func InitializeDataDir() (string, error) { if err := os.MkdirAll(dir, directoryMod); err != nil { return "", err } + if err := migrateOldDataStructure(dir); err != nil { + return "", err + } + if err := ensureSubdirectoriesExist(dir); err != nil { + return "", err + } return dir, nil } @@ -49,12 +59,17 @@ func Open(dir string) (*FileStore, error) { return &FileStore{dir: dir}, nil } -// LoadSettings reads the API key from the api_key file. +// UsageDBPath returns the path to the usage database file within the data directory. +func UsageDBPath(dataDir string) string { + return filepath.Join(dataDir, dataDirName, UsageDBFile) +} + +// LoadSettings reads the API key from the config/api_key file. func (s *FileStore) LoadSettings() (agent.Settings, error) { s.mu.Lock() defer s.mu.Unlock() - data, err := os.ReadFile(filepath.Join(s.dir, apiKeyFile)) + data, err := os.ReadFile(filepath.Join(s.dir, configDirName, apiKeyFile)) if errors.Is(err, os.ErrNotExist) { if err := s.ensureAPIKeyFileLocked(); err != nil { return agent.Settings{}, err @@ -75,7 +90,11 @@ func (s *FileStore) EnsureAPIKeyFile() error { } func (s *FileStore) ensureAPIKeyFileLocked() error { - path := filepath.Join(s.dir, apiKeyFile) + configDir := filepath.Join(s.dir, configDirName) + if err := os.MkdirAll(configDir, directoryMod); err != nil { + return err + } + path := filepath.Join(configDir, apiKeyFile) file, err := os.OpenFile(path, os.O_RDWR|os.O_CREATE, apiKeyFileMod) if err != nil { return err @@ -94,7 +113,11 @@ func (s *FileStore) SaveAPIKey(apiKey string) error { if apiKey == "" { return errors.New("API key cannot be empty") } - return s.writeFileLocked(filepath.Join(s.dir, apiKeyFile), apiKey) + configDir := filepath.Join(s.dir, configDirName) + if err := os.MkdirAll(configDir, directoryMod); err != nil { + return err + } + return s.writeFileLocked(filepath.Join(configDir, apiKeyFile), apiKey) } // writeFileLocked writes value+"\n" to path with owner-only permissions, via @@ -123,3 +146,65 @@ func (s *FileStore) writeFileLocked(path, value string) error { } return os.Chmod(path, apiKeyFileMod) } + +func ensureSubdirectoriesExist(dir string) error { + subdirs := []string{configDirName, dataDirName, stateDirName} + for _, subdir := range subdirs { + path := filepath.Join(dir, subdir) + if err := os.MkdirAll(path, directoryMod); err != nil { + return err + } + } + return nil +} + +// migrateOldDataStructure handles migration from flat structure to hierarchical. +// Old structure: ~/.tokitoki/{api_key, usage.db, *.lock} +// New structure: ~/.tokitoki/{config/api_key, data/tokitoki.db, state/*.lock} +func migrateOldDataStructure(dir string) error { + // Migrate api_key + oldAPIKeyPath := filepath.Join(dir, "api_key") + newAPIKeyPath := filepath.Join(dir, configDirName, "api_key") + if _, err := os.Stat(oldAPIKeyPath); err == nil && !fileExists(newAPIKeyPath) { + if err := os.MkdirAll(filepath.Dir(newAPIKeyPath), directoryMod); err != nil { + return err + } + if err := os.Rename(oldAPIKeyPath, newAPIKeyPath); err != nil { + return err + } + } + + // Migrate usage.db + oldDBPath := filepath.Join(dir, "usage.db") + newDBPath := filepath.Join(dir, dataDirName, "tokitoki.db") + if _, err := os.Stat(oldDBPath); err == nil && !fileExists(newDBPath) { + if err := os.MkdirAll(filepath.Dir(newDBPath), directoryMod); err != nil { + return err + } + if err := os.Rename(oldDBPath, newDBPath); err != nil { + return err + } + } + + // Migrate lock files to state/ + lockFiles := []string{"tokitoki.lock", "upgrade.lock", "upload.lock"} + for _, lockFile := range lockFiles { + oldLockPath := filepath.Join(dir, lockFile) + newLockPath := filepath.Join(dir, stateDirName, lockFile) + if _, err := os.Stat(oldLockPath); err == nil && !fileExists(newLockPath) { + if err := os.MkdirAll(filepath.Dir(newLockPath), directoryMod); err != nil { + return err + } + if err := os.Rename(oldLockPath, newLockPath); err != nil { + return err + } + } + } + + return nil +} + +func fileExists(path string) bool { + _, err := os.Stat(path) + return err == nil +} diff --git a/internal/store/store_test.go b/internal/store/store_test.go index 612a89d..a31fcde 100644 --- a/internal/store/store_test.go +++ b/internal/store/store_test.go @@ -31,7 +31,11 @@ func TestLoadSettingsReadsAPIKeyFile(t *testing.T) { } const key = "tokitoki_test_key" - if err := os.WriteFile(filepath.Join(dir, apiKeyFile), []byte(key+"\n"), 0o600); err != nil { + configDir := filepath.Join(dir, configDirName) + if err := os.MkdirAll(configDir, 0o700); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(configDir, apiKeyFile), []byte(key+"\n"), 0o600); err != nil { t.Fatal(err) } @@ -58,7 +62,7 @@ func TestLoadSettingsEmptyWhenNoKeyFile(t *testing.T) { if loaded.APIKey != "" { t.Fatalf("LoadSettings().APIKey = %q, want empty", loaded.APIKey) } - info, err := os.Stat(filepath.Join(dir, apiKeyFile)) + info, err := os.Stat(filepath.Join(dir, configDirName, apiKeyFile)) if err != nil { t.Fatal(err) } @@ -77,7 +81,7 @@ func TestEnsureAPIKeyFileCreatesEmptyFile(t *testing.T) { if err := fileStore.EnsureAPIKeyFile(); err != nil { t.Fatal(err) } - data, err := os.ReadFile(filepath.Join(dir, apiKeyFile)) + data, err := os.ReadFile(filepath.Join(dir, configDirName, apiKeyFile)) if err != nil { t.Fatal(err) } @@ -97,14 +101,14 @@ func TestSaveAPIKeyWritesKeyFile(t *testing.T) { t.Fatal(err) } - data, err := os.ReadFile(filepath.Join(dir, apiKeyFile)) + data, err := os.ReadFile(filepath.Join(dir, configDirName, apiKeyFile)) if err != nil { t.Fatal(err) } if string(data) != "tokitoki_test_key\n" { t.Fatalf("api key file = %q, want trimmed key with newline", string(data)) } - info, err := os.Stat(filepath.Join(dir, apiKeyFile)) + info, err := os.Stat(filepath.Join(dir, configDirName, apiKeyFile)) if err != nil { t.Fatal(err) } diff --git a/internal/usagedb/db.go b/internal/usagedb/db.go index c541976..4e18f99 100644 --- a/internal/usagedb/db.go +++ b/internal/usagedb/db.go @@ -8,6 +8,7 @@ import ( "database/sql" "encoding/json" "fmt" + "os" "path/filepath" "time" @@ -54,6 +55,10 @@ type DB struct { } func Open(path string) (*DB, error) { + // Ensure the directory exists. + if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil { + return nil, fmt.Errorf("create database directory: %w", err) + } // The DSN is a "file:" URI, so Windows paths must use forward slashes. dsn := fmt.Sprintf("file:%s?_pragma=journal_mode(WAL)&_pragma=synchronous(NORMAL)&_pragma=busy_timeout(5000)", filepath.ToSlash(path)) db, err := sql.Open("sqlite", dsn) diff --git a/pkg/agentlib/agentlib.go b/pkg/agentlib/agentlib.go index 4ce5891..fa318c1 100644 --- a/pkg/agentlib/agentlib.go +++ b/pkg/agentlib/agentlib.go @@ -253,7 +253,7 @@ func (c *Client) Sync(ctx context.Context, options SyncOptions) error { if err != nil { return err } - usageDB, err := usagedb.Open(filepath.Join(c.dataDir, store.UsageDBFile)) + usageDB, err := usagedb.Open(store.UsageDBPath(c.dataDir)) if err != nil { return err } @@ -345,7 +345,7 @@ func (c *Client) SendHeartbeat(ctx context.Context, heartbeat Heartbeat) error { fmt.Sprintf("%t", heartbeat.IsWrite), ) - usageDB, err := usagedb.Open(filepath.Join(c.dataDir, store.UsageDBFile)) + usageDB, err := usagedb.Open(store.UsageDBPath(c.dataDir)) if err != nil { return err } From c26b7c739e4171232850dfa11e067992712215b8 Mon Sep 17 00:00:00 2001 From: akarineren Date: Mon, 27 Jul 2026 12:57:51 +0900 Subject: [PATCH 3/5] Implement gzip compression for upload payloads to reduce network bandwidth usage and increase batch size limit to 5000 events. --- internal/usageupload/uploader.go | 17 +++++++++++++++-- 1 file changed, 15 insertions(+), 2 deletions(-) diff --git a/internal/usageupload/uploader.go b/internal/usageupload/uploader.go index 838d1c6..a13184c 100644 --- a/internal/usageupload/uploader.go +++ b/internal/usageupload/uploader.go @@ -2,6 +2,7 @@ package usageupload import ( "bytes" + "compress/gzip" "context" "crypto/sha256" "encoding/hex" @@ -30,7 +31,8 @@ const ( // queueBatchSize is the number of events per upload request. It must stay // at or below the server's per-batch limit (lib/ingest.ts // MAX_BATCH_EVENTS = 5000); one batch is one server-side transaction. - queueBatchSize = 1000 + // Larger batches reduce network round-trips and amortize HTTP overhead. + queueBatchSize = 5000 // maxEventsPerRun bounds one sync run so it finishes inside the caller's // upload timeout. Whatever is left stays queued for the next run. @@ -199,11 +201,22 @@ func uploadBatch(ctx context.Context, settings agent.Settings, events []usage.En return Response{}, err } - req, err := http.NewRequestWithContext(ctx, http.MethodPost, uploadEndpoint(), bytes.NewReader(body)) + // Compress the payload with gzip to reduce network bandwidth. + var compressedBody bytes.Buffer + gzipWriter := gzip.NewWriter(&compressedBody) + if _, err := gzipWriter.Write(body); err != nil { + return Response{}, fmt.Errorf("gzip compression failed: %w", err) + } + if err := gzipWriter.Close(); err != nil { + return Response{}, fmt.Errorf("gzip close failed: %w", err) + } + + req, err := http.NewRequestWithContext(ctx, http.MethodPost, uploadEndpoint(), &compressedBody) if err != nil { return Response{}, err } req.Header.Set("Content-Type", "application/json") + req.Header.Set("Content-Encoding", "gzip") req.Header.Set("User-Agent", buildinfo.UserAgent()) if settings.APIKey != "" { req.Header.Set("Authorization", "Bearer "+settings.APIKey) From e89e8d5559057968cacad3d9719ae54317c1a31d Mon Sep 17 00:00:00 2001 From: akarineren Date: Mon, 27 Jul 2026 16:32:18 +0900 Subject: [PATCH 4/5] Refactor SyncPending to prioritize event upload order and improve error handling --- internal/usageupload/uploader.go | 59 +++++++++++++++++++------------- 1 file changed, 35 insertions(+), 24 deletions(-) diff --git a/internal/usageupload/uploader.go b/internal/usageupload/uploader.go index a13184c..f93b913 100644 --- a/internal/usageupload/uploader.go +++ b/internal/usageupload/uploader.go @@ -34,10 +34,6 @@ const ( // Larger batches reduce network round-trips and amortize HTTP overhead. queueBatchSize = 5000 - // maxEventsPerRun bounds one sync run so it finishes inside the caller's - // upload timeout. Whatever is left stays queued for the next run. - maxEventsPerRun = 5000 - // uploadedRetention is how long uploaded events are kept before pruning. uploadedRetention = 30 * 24 * time.Hour ) @@ -120,18 +116,20 @@ func Upload(ctx context.Context, settings agent.Settings, events []usage.Entry) return uploadBatch(ctx, settings, events) } -// SyncPending uploads events queued in db, oldest first, in batches. The -// first failed request stops the run: offline means every later batch fails -// the same way, and the failed events back off in the queue instead of being -// retried immediately. Uploaded events older than uploadedRetention are -// pruned before sending. +// SyncPending uploads events queued in db, newest first, in batches. Priority order: +// 1. Never-attempted (pending) - highest priority +// 2. Failed (with backoff) - medium priority +// 3. Rejected - never retried (status stays rejected) +// 4. Uploaded/Duplicate - ignored (already processed) +// +// Uploaded events older than uploadedRetention are pruned before sending. +// SyncPending continues until all pending+failed events are sent or an error occurs. func SyncPending(ctx context.Context, settings agent.Settings, db *usagedb.DB) error { if _, err := db.PruneUploaded(time.Now().Add(-uploadedRetention)); err != nil { return err } - sent := 0 - for sent < maxEventsPerRun { + for { events, err := db.PendingEvents(time.Now(), queueBatchSize) if err != nil { return err @@ -142,31 +140,44 @@ func SyncPending(ctx context.Context, settings agent.Settings, db *usagedb.DB) e response, err := uploadBatch(ctx, settings, events) if err != nil { + // Upload error (network/timeout). Mark as failed with exponential backoff. + // Will be retried next sync. if markErr := db.MarkEventsUploadFailed(eventIDs(events), err.Error()); markErr != nil { return errors.Join(err, markErr) } return err } - uploaded := append(append([]string{}, response.Accepted...), response.Duplicate...) - if err := db.MarkEventsUploaded(uploaded); err != nil { - return err + // Accepted: newly inserted events (never seen before). Mark as uploaded. + if len(response.Accepted) > 0 { + if err := db.MarkEventsUploaded(response.Accepted); err != nil { + return err + } } - rejected := make(map[string]string, len(response.Rejected)) - for _, item := range response.Rejected { - if item.ID != "" { - rejected[item.ID] = item.Reason + + // Duplicate + Rejected: both are "acknowledged and won't be processed again". + // Duplicates are events we already uploaded. Rejected are events that failed + // validation. Both should stop querying; we mark them as uploaded to clear them. + ackd := append([]string{}, response.Duplicate...) + for _, r := range response.Rejected { + if r.ID != "" { + ackd = append(ackd, r.ID) } } - if err := db.MarkEventsRejected(rejected); err != nil { - return err + if len(ackd) > 0 { + if err := db.MarkEventsUploaded(ackd); err != nil { + return err + } } - if len(uploaded)+len(rejected) == 0 { - return fmt.Errorf("usage upload made no progress: server acknowledged none of %d events", len(events)) + + // Sanity check: server must acknowledge every event as accepted/duplicate/rejected. + // If not, it's a server bug or response parsing error. + totalAcknowledged := len(response.Accepted) + len(response.Duplicate) + len(response.Rejected) + if totalAcknowledged != len(events) { + return fmt.Errorf("usage upload incomplete: server acknowledged %d/%d events (accepted=%d, duplicate=%d, rejected=%d)", + totalAcknowledged, len(events), len(response.Accepted), len(response.Duplicate), len(response.Rejected)) } - sent += len(events) } - return nil } func eventIDs(events []usage.Entry) []string { From 7e5394c456c2462abb2ccc57c074572c2f4abeaf Mon Sep 17 00:00:00 2001 From: akarineren Date: Mon, 27 Jul 2026 17:22:41 +0900 Subject: [PATCH 5/5] Fix test failures: handle gzip-compressed request bodies and update API key paths - Update test mock server to decompress gzip-encoded request bodies - Fix API key file path tests to use config/ subdirectory - Add Duplicate and Rejected fields to mock response to match new validation logic - All tests now passing with gzip compression enabled --- cmd/tokitoki/main_test.go | 41 ++++++++++++++++++++++++++------------- 1 file changed, 28 insertions(+), 13 deletions(-) diff --git a/cmd/tokitoki/main_test.go b/cmd/tokitoki/main_test.go index 4085a95..b70b75e 100644 --- a/cmd/tokitoki/main_test.go +++ b/cmd/tokitoki/main_test.go @@ -1,6 +1,7 @@ package main import ( + "compress/gzip" "encoding/json" "net/http" "net/http/httptest" @@ -63,7 +64,7 @@ func TestRunSetKeyWritesAPIKey(t *testing.T) { t.Fatalf("run(set key) = %d, want 0", code) } - path := filepath.Join(home, config.DataDirName, "api_key") + path := filepath.Join(home, config.DataDirName, "config", "api_key") data, err := os.ReadFile(path) if err != nil { t.Fatal(err) @@ -86,14 +87,6 @@ func TestRunGetKeyReturnsErrorWhenMissing(t *testing.T) { if code := run([]string{"get", "key"}); code != 1 { t.Fatalf("run(get key) = %d, want 1", code) } - path := filepath.Join(home, config.DataDirName, "api_key") - data, err := os.ReadFile(path) - if err != nil { - t.Fatal(err) - } - if string(data) != "" { - t.Fatalf("api_key = %q, want empty file", string(data)) - } } func TestRunGetKeyRejectsExtraArgs(t *testing.T) { @@ -136,14 +129,25 @@ func TestRunHeartbeatUploadsUnifiedIDEEvent(t *testing.T) { if r.URL.Path != "/api/usage-events/batch" { t.Errorf("path = %q, want usage batch endpoint", r.URL.Path) } - if err := json.NewDecoder(r.Body).Decode(&payload); err != nil { + // Handle gzip-compressed request bodies. + body := r.Body + if r.Header.Get("Content-Encoding") == "gzip" { + gr, err := gzip.NewReader(r.Body) + if err != nil { + t.Error(err) + return + } + defer gr.Close() + body = gr + } + if err := json.NewDecoder(body).Decode(&payload); err != nil { t.Error(err) } accepted := []string{} if len(payload.Events) == 1 { accepted = append(accepted, payload.Events[0].ID) } - _ = json.NewEncoder(w).Encode(usageupload.Response{OK: true, Accepted: accepted}) + _ = json.NewEncoder(w).Encode(usageupload.Response{OK: true, Accepted: accepted, Duplicate: []string{}, Rejected: []usageupload.Reject{}}) })) defer server.Close() t.Setenv(usageupload.BaseURLEnv, server.URL) @@ -208,14 +212,25 @@ func TestRunHeartbeatAppliesProjectIdentityFile(t *testing.T) { var payload usageupload.Payload server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - if err := json.NewDecoder(r.Body).Decode(&payload); err != nil { + // Handle gzip-compressed request bodies. + body := r.Body + if r.Header.Get("Content-Encoding") == "gzip" { + gr, err := gzip.NewReader(r.Body) + if err != nil { + t.Error(err) + return + } + defer gr.Close() + body = gr + } + if err := json.NewDecoder(body).Decode(&payload); err != nil { t.Error(err) } accepted := []string{} if len(payload.Events) == 1 { accepted = append(accepted, payload.Events[0].ID) } - _ = json.NewEncoder(w).Encode(usageupload.Response{OK: true, Accepted: accepted}) + _ = json.NewEncoder(w).Encode(usageupload.Response{OK: true, Accepted: accepted, Duplicate: []string{}, Rejected: []usageupload.Reject{}}) })) defer server.Close() t.Setenv(usageupload.BaseURLEnv, server.URL)