collect

package
v0.1.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 3, 2026 License: MIT Imports: 7 Imported by: 0

Documentation

Overview

Package collect drives one or more collection cycles: it iterates the adapter registry, reads each discovered source read-only, and persists the observations into the append-only store. EventLevel observations are stored directly (deduplicated on DedupKey). Aggregate observations are turned into synthetic immutable events via a monotonic-with-reset delta against the last stored accumulator state, so a source that later shrinks (compaction, deletion, reset) can never reduce a previously-reported total.

TWO ENTRY POINTS, AND NOTHING AROUND THEM: RunOnce performs one pass and Run performs one every interval until its context is cancelled (issue #72, decision 3). Everything a MACHINE needs around a long-running collector - a pidfile and its lock, a recorded build identity, a stop-and-wait, a watch on the executable for an upgrade - is deliberately absent, because those are this project's answers to questions a consumer embedding the package has already answered its own way. aiusage's own daemon keeps all of it in internal/daemon, which composes this package exactly the way any other consumer would.

The store is taken as Store, an interface this package declares over the six methods a pass actually uses (decision 8), so *store.Ledger satisfies it without either package naming the other's shape.

INJECTION SEAMS: RunOnce accepts a Pricer interface to price each new event at ingest time, and an optional Refresher to update the price table before stamping. A versioned pricer also fills missing historical costs when the store supports price sync. A Pricer returns microUSD, source, and ok; ok=false leaves the event unpriced (CostMicroUSD = nil). Cost stamped by an adapter (attached to the event before collection) is never overwritten by the ladder: the Pricer is consulted only for events that arrive UNPRICED. WithPricer() / WithoutRaw() are Options that configure a cycle; Options are composable and independent.

INCREMENTAL COLLECTION: adapters that implement adapter.Incremental can supply a checkpoint (SourceCheckpoint) to skip unchanged sources — a seam for efficiency, not correctness. The checkpoint rides the same transaction as its events, so a crash cannot advance a checkpoint past the rows it accounts for. A checkpoint load failure stops that source: some adapters keep monetary baselines in the checkpoint, so a full read could double charge history.

SINGLE-TRANSACTION OBSERVATION: events, activity records, and turn contexts from one source observation land in ONE transaction per source, alongside the source checkpoint (for append-only sources; for aggregate snapshots the checkpoint lands on the final snapshot). A crashed or collided cell re-reads the whole source. This transaction boundary is the unit of consistency: an observer cannot see one part of an observation without the others.

ACTIVITY and TURN CONTEXT: The collector records which tool was called (activity_events) and which subagent/skill/MCP tool/server/plugin each turn ran under (usage_turn_context). Activity rows carry no cost column — cost is derived on read by joining the usage event the activity names, so one turn's tokens are not multiplied by the number of calls it made. Turn contexts carry no token information — they are a property of the turn (dimension attribute), and a turn commonly carries several (subagent AND skill AND MCP tool, etc.), so a dimension-blind query would overstate cost by summing the same turn multiple times.

Index

Examples

Constants

This section is empty.

Variables

This section is empty.

Functions

func Run

func Run(ctx context.Context, interval time.Duration, reg *adapter.Registry, st Store, dc adapter.DiscoverConfig, opts ...Option) error

Run collects on a ticker until ctx is cancelled: one pass immediately, then one every interval. It is the whole of the long-running half of this package (issue #72, decision 3) - a loop over RunOnce and nothing else.

WHAT IT DELIBERATELY DOES NOT DO: take a pidfile lock, write a pid, record a build identity, watch its own executable, or install a signal handler. Those are how one machine supervises one process, and a consumer embedding this package has its own answers to all of them - a supervisor, a container, a goroutine inside a larger service. Baking this project's answers in would make every one of those callers fight them. aiusage's own daemon keeps them, one level down in internal/daemon, composing this package the way any other consumer would.

The consequence a caller owns: NOTHING here prevents two Runs against one database. They would not corrupt it - the ledger is append-only and dedup keys collide harmlessly - but two passes reading one aggregate accumulator at the same time can each derive the same delta under a different dedup key and append it twice. One Run per database.

A per-pass failure is not fatal, because RunOnce's own errors are per-source and it is expected to be run on a ticker: the loop reports each pass through WithCycleCallback (if one was given) and keeps going. Run returns nil on cancellation - being asked to stop is not a failure - so an error from it means the loop could not continue at all.

Types

type CycleStats

type CycleStats struct {
	Adapters       int      // adapters iterated
	Sources        int      // sources discovered + collected
	SourcesFailed  int      // sources that recorded errors and stored nothing
	EventsInserted int      // new dedup keys actually written (event + synthetic)
	EventsSeen     int      // observed event-level records (pre-dedup)
	Snapshots      int      // aggregate snapshots observed
	Errors         []string // non-fatal per-adapter / per-source errors
	PricesSynced   int      // existing unpriced rows given a computed cost

	// ActivitySeen counts observed tool/skill/hook invocations (pre-dedup) and
	// ActivityInserted the new activity dedup keys actually written. They are
	// kept apart from the event counts because they measure a different ledger:
	// a cycle that inserts no events can still insert activity, and reading one
	// as the other would misreport both.
	ActivitySeen     int
	ActivityInserted int

	// TurnContextsSeen counts observed (turn, dimension) attributions — which
	// subagent, skill, MCP tool, MCP server or plugin each turn ran under — and
	// TurnContextsInserted the new ones actually written. A context is one per
	// (USAGE EVENT, DIMENSION), so ONE turn can contribute up to five: these
	// count rows, not turns, and they track the event counts rather than the
	// activity ones, being the turns a context spent and not the calls it made.
	TurnContextsSeen     int
	TurnContextsInserted int

	// RollupRebuilt reports that the pass found the derived rollup out
	// of step with the ledger and rebuilt it before collecting. Expected once,
	// on the first pass after the v4 migration; anywhere else it means a write
	// reached the ledger without its rollup delta and is worth noticing.
	RollupRebuilt bool

	// Canceled reports that the context was cancelled mid-pass and the cycle
	// stopped early. Every count above is then partial — in particular the
	// adapter and source being processed are already counted — so a truncated
	// cycle must never be read (or logged) as a completed one.
	Canceled bool

	// CodeChangesUpdated counts per-turn snapshots inserted or changed.
	CodeChangesUpdated int
}

CycleStats reports the outcome of a single RunOnce.

func RunOnce

func RunOnce(ctx context.Context, reg *adapter.Registry, st Store, dc adapter.DiscoverConfig, opts ...Option) (CycleStats, error)

RunOnce performs one full collection pass. Per-source and per-adapter errors are non-fatal: they are appended to CycleStats.Errors and collection continues. RunOnce only returns a non-nil error for failures that prevent the cycle from making any meaningful progress (none currently — the loop is fully resilient), so callers may safely run it on a ticker.

A cancelled context truncates the pass: RunOnce returns ctx.Err() together with CycleStats.Canceled set, and the counts in those stats cover only the work done so far. Nothing is lost by the truncation — checkpoints ride the event transactions — but the stats must be reported as partial.

Example
package main

import (
	"context"
	"fmt"
	"os"
	"path/filepath"

	"github.com/RandomCodeSpace/aiusage-core/adapter"
	"github.com/RandomCodeSpace/aiusage-core/adapter/codex"
	"github.com/RandomCodeSpace/aiusage-core/collect"
	"github.com/RandomCodeSpace/aiusage-core/pricing"
	"github.com/RandomCodeSpace/aiusage-core/store"
)

func main() {
	// This example supplies a fixture root. An application supplies the root
	// of its existing source and keeps its output ledger between passes.
	dir, err := os.MkdirTemp("", "aiusage-core-example-")
	if err != nil {
		panic(err)
	}
	defer os.RemoveAll(dir)
	root := filepath.Join(dir, "codex")
	if err := os.MkdirAll(filepath.Join(root, "sessions"), 0o700); err != nil {
		panic(err)
	}
	fixture := `{"type":"turn_context","payload":{"model":"gpt-5-codex"}}
{"type":"event_msg","timestamp":"2026-05-29T10:00:00Z","payload":{"type":"token_count","info":{"last_token_usage":{"input_tokens":1000,"cached_input_tokens":400,"output_tokens":200,"reasoning_output_tokens":50,"total_tokens":1200}}}}
`
	if err := os.WriteFile(filepath.Join(root, "sessions", "example.jsonl"), []byte(fixture), 0o600); err != nil {
		panic(err)
	}
	ledger, err := store.Open(filepath.Join(dir, "usage.db"))
	if err != nil {
		panic(err)
	}
	defer ledger.Close()

	ctx := context.Background()
	registry := adapter.NewRegistry(codex.New())
	discovery := adapter.DiscoverConfig{Overrides: map[string]string{model.ToolCodex: root}}
	prices := pricing.New(pricing.Options{})
	stats, err := collect.RunOnce(ctx, registry, ledger, discovery,
		collect.WithPricer(prices), collect.WithoutRaw())
	if err != nil || len(stats.Errors) != 0 {
		panic(fmt.Sprintf("collection failed: %v %v", err, stats.Errors))
	}
	summary, err := ledger.Summarize(ctx, store.Filter{})
	if err != nil {
		panic(err)
	}
	fmt.Printf("inserted=%d tokens=%d\n", stats.EventsInserted, summary.Totals.Total)
}
Output:
inserted=1 tokens=1200

func (CycleStats) AllFailed

func (s CycleStats) AllFailed() bool

AllFailed reports whether the cycle produced only errors: every discovered source failed outright, or discovery itself yielded nothing but errors. `once` maps this to its nonzero exit for cron use; partial failure stays 0.

type Option

type Option func(*cycleOptions)

Option configures a collection cycle.

func WithCycleCallback

func WithCycleCallback(fn func(CycleStats, error)) Option

WithCycleCallback reports the outcome of each pass Run makes: the stats and the error RunOnce returned, in that order, on Run's own goroutine.

It exists because Run is a loop with no return value until it stops, so without it a long-running collector is a black box - a consumer could not log a pass, export a counter, or notice that every source has been failing for an hour. RunOnce ignores it: a caller holding the stats already has them.

The callback runs BETWEEN passes and blocks the next one, which is the honest ordering (a pass is not reported until it is over) and makes a slow callback a slow collector. Keep it to logging or a metric.

func WithPricer

func WithPricer(p Pricer) Option

WithPricer stamps a cost on every new event from the price table in effect at ingest. A pricer with Revision() string also enables historical price sync on stores that support it. Without a pricer, events are stored unpriced.

func WithoutRaw

func WithoutRaw() Option

WithoutRaw drops every adapter's raw audit payload before it reaches the store (config privacy.no_raw). The collector is the choke point every adapter's output passes through, so the switch holds for adapters that never learn about it — including ones added later — and covers both raw columns: usage_events.raw via the events, aggregate_state.raw and the synthetic delta event via the snapshots.

type Pricer

type Pricer interface {
	PriceEvent(e model.UsageEvent) (microUSD int64, source string, ok bool)
}

Pricer stamps a cost on a usage event just before it is appended. It returns the cost in micro-USD and the price_source that produced it. ok=false leaves the event unpriced (a NULL cost column) — the honest state, since a stored 0 would claim the request was free. Reporting prices NULL rows from the current table at display time instead. A versioned pricer also lets collection fill missing stored costs when rates become available later.

type Refresher

type Refresher interface {
	Refresh(ctx context.Context) error
}

Refresher is the optional half of Pricer: a price ladder that can update itself from upstream. RunOnce offers it one chance per cycle; the implementation is expected to throttle itself and to fail silently.

type Store

type Store interface {
	EnsureRollup(ctx context.Context) (bool, error)
	Checkpoint(ctx context.Context, tool, sourcePath string) (*model.SourceCheckpoint, error)
	ApplyBatch(ctx context.Context, b store.ObservationBatch) (store.Applied, error)
	ApplyEvents(ctx context.Context, events []model.UsageEvent, cp *model.SourceCheckpoint) (int, error)
	LastState(ctx context.Context, tool, key string) (*model.AggregateSnapshot, error)
	ApplySnapshot(ctx context.Context, events []model.UsageEvent, state model.AggregateSnapshot, cp *model.SourceCheckpoint) (int, error)
}

Store is what a collection pass needs from a ledger, and nothing else: six methods, all of which a *store.Ledger already satisfies.

It is declared HERE, at the point of use, rather than exported by package store (issue #72, decision 8). A fat interface beside the implementation has to name every method the implementation has, so it grows whenever the store grows and every fake in the module breaks on a method the fake's own test never calls. Declared here it names what THIS package consumes: adding a query to the store cannot break the collector's tests, and a caller wiring a different ledger in has six methods to write instead of thirty.

The contracts these carry are the store's, not this package's, and they are not restated below - see store.Ledger for the single-transaction and idempotence rules each one promises. What matters here is which of them the collector depends on:

  • EnsureRollup, once per pass, before anything is appended.
  • Checkpoint, per source, for the incremental gate.
  • ApplyBatch, per source observation: the one transaction that carries events, activity, turn contexts and the checkpoint together.
  • LastState / ApplySnapshot, per aggregate cell, for the monotonic-with-reset delta and its baseline.
  • ApplyEvents, for the checkpoint-only write of an unchanged aggregate cell.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL