csfpg

package
v0.2.3 Latest Latest
Warning

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

Go to latest
Published: Oct 3, 2026 License: Apache-2.0 Imports: 14 Imported by: 0

README

csfpg: CSF's PostgreSQL backend

candace/ipc/db/csfpg is CSF's PostgreSQL backend: the only place a CSF process opens a PostgreSQL pool, the owner of CSF's schema and its migrations, and the home of the queries sqlc generates against that schema. Migrations are the only schema source. It follows the ipc/db rule: one package per database backend CSF reaches, named csf<backend>; the stores that use it live in services/.

Exported Role
IDB The capability a service receives: sqlc's DBTX plus Begin, BeginTx and Ping.
Settings, SettingsFromEnvironment What the binary opens, read from its own configuration.
OpenPool, Pool, Pool.OpenSQL The binary opens one pool, hands it out as an IDB, and closes it after its services stop.
ApplySchema, SchemaVersionTable CSF's one migration, its schema, applied under csf_schema_version.
New, Queries and the Csf* models sqlc-generated from queries.sql, jobs.sql, cron.sql and agent_configurations.sql against the schema; a service passes its IDB to New. The query contract is in QUERIES.md.
mocks.MockIDB, mocks.MockTx Generated gomock doubles (go generate ./ipc/db/csfpg) for specs.

Schema

CSF ships exactly one migration, schema/001_init.sql. It assumes a clean database and is edited in place while CSF's schema evolves. candace/pkg/sqlmigrate applies it and records it in CSF's own version table, csf_schema_version; sqlc.yaml generates this package's queries from the same file (see QUERIES.md).

An application that embeds CSF keeps its own migrations, numbered independently and recorded in its own version table. It may use any migration tool for them. With sqlmigrate, the binary makes two calls on one database:

handle := pool.OpenSQL()
defer handle.Close()
// CSF's schema, under csf_schema_version.
if err := csfpg.ApplySchema(ctx, handle); err != nil {
	return err
}
// The application's own migrations, under its own version table.
if err := sqlmigrate.ApplyVersioned(ctx, handle, migrations, "migrations", "myapp_schema_version"); err != nil {
	return err
}

candace/examples/csfpg-consumer/migrations/001_example_notes.sql is a copyable starting point. schema_test.go loads that exact directory and stacks it on CSF's schema in one database, so the example cannot rot.

Tests

Specs never open a real database. A unit of code that uses the database takes csfpg.IDB, and its specs use mocks.MockIDB. Specs that need the schema itself run it on pgmem, loading the real 001_init.sql with no DDL in Go. The generated queries run on pgmem too: pgmem.Schema.OpenPGX serves the pgx surface and satisfies IDB, so queries_test.go passes it to New and runs one query file per spec.

The labelled exception is the opt-in acceptance tier: specs behind the acceptance build tag that run against a disposable PostgreSQL named by CANDACE_CSF_TEST_DATABASE_URL (the deploy service's specs use CANDACEOS_STORE_TEST_DATABASE_URL), opening pools only through csfpg.OpenPool. It holds only what pgmem cannot yet prove: this package's pool and ApplySchema on real PostgreSQL, and the deploy service's suites under services/deploy, whose SQL uses PL/pgSQL triggers, LEFT JOIN LATERAL, arrays and DISTINCT ON views.

go test -race -tags acceptance ./ipc/db/csfpg/ ./services/deploy/... ./app/deploy/...

Documentation

Overview

Package csfpg is CSF's PostgreSQL backend: the one owner of database connection pools in a CSF process, CSF's schema and its migrations, and the queries sqlc generates against that schema. It is one of the ipc/db packages, one per database backend CSF reaches, each named csf<backend>; the stores that use it live with their services.

A connection to PostgreSQL crosses the kernel I/O tier, so it is opened only here, by the binary that owns the process. The binary calls OpenPool once with Settings it read from its own configuration, hands the resulting Pool to each service as an IDB, and closes the pool after those services stop. A service receives the IDB in its constructor and never opens or closes a pool itself; pgxpool.New and pgx.Connect appear nowhere else.

Index

Constants

View Source
const SchemaDirectory = "schema"

SchemaDirectory is the directory of Schema that holds the migrations.

View Source
const SchemaVersionTable = "csf_schema_version"

SchemaVersionTable records which of CSF's migrations a database holds. An application that embeds CSF keeps its own migrations in its own version table, applied after CSF's with any tool; the two are numbered independently and never share a table.

Variables

View Source
var ErrMissingURL = errors.New("ipc/db/csfpg: a database URL is required")

ErrMissingURL means the settings name no database.

View Source
var Schema embed.FS

Schema holds CSF's migrations. CSF ships one, schema/001_init.sql: it assumes a clean database and is edited in place while the schema evolves. sqlc generates this package's queries from the same file.

Functions

func ApplySchema

func ApplySchema(ctx context.Context, database *sql.DB) error

ApplySchema brings database up to CSF's schema. It is idempotent: a file the version table records is never applied again. Pass a pool's OpenSQL handle.

func Time

func Time(value pgtype.Timestamptz) time.Time

Time is the Go instant of a pgx timestamp: NULL is the zero time, and a set value is read in UTC.

func Timestamp

func Timestamp(value time.Time) pgtype.Timestamptz

Timestamp is the pgx value of a Go instant: the zero time is NULL, and a set time is UTC at PostgreSQL's microsecond precision, so what a service compares in Go is what the database stored.

func Transact

func Transact(ctx context.Context, database IDB, work func(tx pgx.Tx) error) error

Transact runs work inside one transaction of database and commits it when work returns nil; any error rolls the transaction back and is returned. It is the one transaction shape every store over an IDB uses.

Types

type AcquireExpiredCronOccurrenceParams

type AcquireExpiredCronOccurrenceParams struct {
	LeaseOwner   *string
	LeaseToken   *string
	LeaseUntil   pgtype.Timestamptz
	ClaimedAt    pgtype.Timestamptz
	OccurrenceID string
}

type AdvanceCronTriggerParams

type AdvanceCronTriggerParams struct {
	NextRunAt   pgtype.Timestamptz
	UpdatedAt   pgtype.Timestamptz
	TriggerName string
}

type ClaimJobSubmissionParams

type ClaimJobSubmissionParams struct {
	Executor string
	Kinds    []string
}

type ClaimJobTraceDeliveryParams

type ClaimJobTraceDeliveryParams struct {
	Destination string
	JobID       string
	TraceID     string
}

type ClaimProjectionTaskParams

type ClaimProjectionTaskParams struct {
	LeaseSeconds int32
	MaxAttempts  int32
}

type ClaimProjectionTaskRow

type ClaimProjectionTaskRow struct {
	SourceID        string
	Revision        string
	Status          CsfProjectionStatus
	Attempts        int32
	LeaseGeneration int64
	LeaseUntil      pgtype.Timestamptz
	NextAttemptAt   pgtype.Timestamptz
	LastError       string
	CreatedAt       pgtype.Timestamptz
	UpdatedAt       pgtype.Timestamptz
}

type CompleteProjectionTaskParams

type CompleteProjectionTaskParams struct {
	SourceID        string
	Revision        string
	LeaseGeneration int64
}

type CountJobStatesRow

type CountJobStatesRow struct {
	Kind     string
	Executor string
	State    CsfJobState
	Jobs     int64
}

type CountProjectionTasksRow

type CountProjectionTasksRow struct {
	Status CsfProjectionStatus
	Count  int64
}

type CreateAgentConfigurationParams

type CreateAgentConfigurationParams struct {
	AgentID                        string
	LangfuseEndpointUrl            string
	LangfusePublicKeySecretRef     string
	LangfuseSecretKeySecretRef     string
	OpensearchEndpointUrl          string
	OpensearchIndex                string
	OpensearchEmbeddingModel       string
	OpensearchCredentialsSecretRef string
}

type CreateDocumentParams

type CreateDocumentParams struct {
	ContentHash string
	ByteSize    int64
	ArtifactRef string
}

type CreateEdgeParams

type CreateEdgeParams struct {
	Relation   string
	AuthorKind string
	AuthorRef  string
	Rationale  string
	FromNodeID string
	ToNodeID   string
}

type CreateNodeParams

type CreateNodeParams struct {
	NodeID              string
	Kind                string
	SymbolKey           string
	Title               string
	Statement           string
	ParentNodeID        *string
	AuthorKind          string
	AuthorRef           string
	CitationContentHash *string
	CitationSourceID    *string
	CitationRevision    *string
	CitationChunkID     *string
	CitationStartByte   *int64
	CitationEndByte     *int64
	TextProjectionRef   *string
	VectorProjectionRef *string
}

type CreateSourceRevisionParams

type CreateSourceRevisionParams struct {
	SourceID             string
	Revision             string
	ContentHash          string
	RawSourceContentHash *string
	SourceUri            string
	Title                string
	MediaType            string
	License              string
	RetrievedAt          pgtype.Timestamptz
}

type CsfAgentConfiguration

type CsfAgentConfiguration struct {
	AgentID                        string
	Revision                       int64
	LangfuseEndpointUrl            string
	LangfusePublicKeySecretRef     string
	LangfuseSecretKeySecretRef     string
	OpensearchEndpointUrl          string
	OpensearchIndex                string
	OpensearchEmbeddingModel       string
	OpensearchCredentialsSecretRef string
	UpdatedAt                      pgtype.Timestamptz
}

type CsfCronOccurrence

type CsfCronOccurrence struct {
	OccurrenceID string
	TriggerName  string
	ScheduledAt  pgtype.Timestamptz
	Status       string
	Attempt      int32
	LeaseOwner   *string
	LeaseToken   *string
	LeaseUntil   pgtype.Timestamptz
	StartedAt    pgtype.Timestamptz
	FinishedAt   pgtype.Timestamptz
	ErrorSummary *string
	SkipReason   *string
	CreatedAt    pgtype.Timestamptz
	UpdatedAt    pgtype.Timestamptz
}

type CsfCronTrigger

type CsfCronTrigger struct {
	TriggerName         string
	ScheduleKind        string
	LocalHour           *int16
	LocalMinute         *int16
	Weekday             *int16
	MonthDay            *int16
	IntervalNanoseconds *int64
	RawExpression       *string
	Timezone            string
	IntervalAnchorAt    pgtype.Timestamptz
	NextRunAt           pgtype.Timestamptz
	CatchUpPolicy       string
	OverlapPolicy       string
	Enabled             bool
	CreatedAt           pgtype.Timestamptz
	UpdatedAt           pgtype.Timestamptz
}

type CsfDocument

type CsfDocument struct {
	ContentHash string
	ByteSize    int64
	ArtifactRef string
	CreatedAt   pgtype.Timestamptz
}

type CsfEdge

type CsfEdge struct {
	FromNodeID string
	ToNodeID   string
	Relation   string
	AuthorKind string
	AuthorRef  string
	Rationale  string
	CreatedAt  pgtype.Timestamptz
}

type CsfIntent

type CsfIntent struct {
	IntentID  string
	SliceID   *string
	Intent    []byte
	CreatedAt pgtype.Timestamptz
}

type CsfJob

type CsfJob struct {
	JobID                 string
	Kind                  string
	Executor              string
	BudgetAccount         string
	Request               []byte
	RequestSha256         string
	State                 CsfJobState
	TotalUnits            int64
	CompletedUnits        int64
	Managed               bool
	ExecutorTarget        string
	ExecutorImage         string
	ExternalID            string
	ArtifactUri           string
	ReservationUsdMicros  int64
	TimeoutSeconds        int32
	Reason                string
	LogStream             string
	LogCursor             string
	InspectionError       string
	CancellationRequested bool
	CleanupConfirmed      bool
	TraceUrl              string
	TraceExportError      string
	TraceRetryAt          pgtype.Timestamptz
	LogDocumentID         string
	LogProjectionError    string
	LogIndexedAt          pgtype.Timestamptz
	LogRetryAt            pgtype.Timestamptz
	CreatedAt             pgtype.Timestamptz
	UpdatedAt             pgtype.Timestamptz
}

type CsfJobBudget

type CsfJobBudget struct {
	Account           string
	LimitUsdMicros    int64
	ReservedUsdMicros int64
	CreatedAt         pgtype.Timestamptz
}

type CsfJobMeasurement

type CsfJobMeasurement struct {
	JobID        string
	Metric       string
	Step         int64
	Value        float64
	RecordedAt   pgtype.Timestamptz
	EvidenceHash string
}

type CsfJobMetricDefinition

type CsfJobMetricDefinition struct {
	JobID       string
	Name        string
	Unit        string
	Description string
}

type CsfJobState

type CsfJobState string
const (
	CsfJobStatePending           CsfJobState = "pending"
	CsfJobStateSubmitting        CsfJobState = "submitting"
	CsfJobStateQueued            CsfJobState = "queued"
	CsfJobStateRunning           CsfJobState = "running"
	CsfJobStateSucceeded         CsfJobState = "succeeded"
	CsfJobStateFailed            CsfJobState = "failed"
	CsfJobStateCancelling        CsfJobState = "cancelling"
	CsfJobStateCancelled         CsfJobState = "cancelled"
	CsfJobStateSubmissionUnknown CsfJobState = "submission_unknown"
)

func (*CsfJobState) Scan

func (e *CsfJobState) Scan(src interface{}) error

type CsfJobTraceDelivery

type CsfJobTraceDelivery struct {
	Destination string
	TraceID     string
	JobID       string
	State       CsfTraceDeliveryState
	Error       string
	CreatedAt   pgtype.Timestamptz
	FinishedAt  pgtype.Timestamptz
}

type CsfNode

type CsfNode struct {
	NodeID              string
	Kind                string
	SymbolKey           string
	Title               string
	Statement           string
	ParentNodeID        *string
	AuthorKind          string
	AuthorRef           string
	CitationContentHash *string
	CitationSourceID    *string
	CitationRevision    *string
	CitationChunkID     *string
	CitationStartByte   *int64
	CitationEndByte     *int64
	TextProjectionRef   *string
	VectorProjectionRef *string
	CreatedAt           pgtype.Timestamptz
}

type CsfProjectionStatus

type CsfProjectionStatus string
const (
	CsfProjectionStatusPending   CsfProjectionStatus = "pending"
	CsfProjectionStatusRunning   CsfProjectionStatus = "running"
	CsfProjectionStatusSucceeded CsfProjectionStatus = "succeeded"
	CsfProjectionStatusFailed    CsfProjectionStatus = "failed"
)

func (*CsfProjectionStatus) Scan

func (e *CsfProjectionStatus) Scan(src interface{}) error

type CsfProjectionTask

type CsfProjectionTask struct {
	SourceID        string
	Revision        string
	Status          CsfProjectionStatus
	Attempts        int32
	LeaseGeneration int64
	LeaseUntil      pgtype.Timestamptz
	NextAttemptAt   pgtype.Timestamptz
	LastError       string
	CreatedAt       pgtype.Timestamptz
	UpdatedAt       pgtype.Timestamptz
}

type CsfSlice

type CsfSlice struct {
	SliceID        string
	Sequence       int64
	Title          string
	Recipe         []byte
	TouchSet       []byte
	Provenance     []byte
	State          CsfSliceState
	AssignmentID   string
	PullRequestUrl string
	Attempts       int32
	Checkpoint     string
	Error          string
	CreatedAt      pgtype.Timestamptz
	UpdatedAt      pgtype.Timestamptz
}

type CsfSliceEdge

type CsfSliceEdge struct {
	FromSliceID string
	ToSliceID   string
	Relation    string
}

type CsfSliceState

type CsfSliceState string
const (
	CsfSliceStateQueued    CsfSliceState = "queued"
	CsfSliceStateRunning   CsfSliceState = "running"
	CsfSliceStatePreempted CsfSliceState = "preempted"
	CsfSliceStateMerged    CsfSliceState = "merged"
	CsfSliceStateFailed    CsfSliceState = "failed"
	CsfSliceStateCanceled  CsfSliceState = "canceled"
)

func (*CsfSliceState) Scan

func (e *CsfSliceState) Scan(src interface{}) error

type CsfSourceRevision

type CsfSourceRevision struct {
	SourceID             string
	Revision             string
	ContentHash          string
	RawSourceContentHash *string
	SourceUri            string
	Title                string
	MediaType            string
	License              string
	RetrievedAt          pgtype.Timestamptz
}

type CsfTraceDeliveryState

type CsfTraceDeliveryState string
const (
	CsfTraceDeliveryStateAttempted CsfTraceDeliveryState = "attempted"
	CsfTraceDeliveryStateSucceeded CsfTraceDeliveryState = "succeeded"
	CsfTraceDeliveryStateAmbiguous CsfTraceDeliveryState = "ambiguous"
)

func (*CsfTraceDeliveryState) Scan

func (e *CsfTraceDeliveryState) Scan(src interface{}) error

type DBTX

type DBTX interface {
	Exec(context.Context, string, ...interface{}) (pgconn.CommandTag, error)
	Query(context.Context, string, ...interface{}) (pgx.Rows, error)
	QueryRow(context.Context, string, ...interface{}) pgx.Row
}

type DisableCronTriggerParams

type DisableCronTriggerParams struct {
	UpdatedAt   pgtype.Timestamptz
	TriggerName string
}

type EnqueueProjectionTaskParams

type EnqueueProjectionTaskParams struct {
	SourceID string
	Revision string
}

type EnsureJobBudgetParams

type EnsureJobBudgetParams struct {
	Account        string
	LimitUsdMicros int64
}

type FailProjectionTaskParams

type FailProjectionTaskParams struct {
	MaxAttempts      int32
	RetryMaxSeconds  int32
	RetryBaseSeconds int32
	LastError        string
	SourceID         string
	Revision         string
	LeaseGeneration  int64
}

type FinishCronOccurrenceParams

type FinishCronOccurrenceParams struct {
	Status       string
	FinishedAt   pgtype.Timestamptz
	ErrorSummary *string
	OccurrenceID string
	LeaseToken   *string
}

type FinishJobSubmissionParams

type FinishJobSubmissionParams struct {
	JobID      string
	State      CsfJobState
	ExternalID string
	Reason     string
}

type FinishJobTraceDeliveryParams

type FinishJobTraceDeliveryParams struct {
	State       CsfTraceDeliveryState
	Error       string
	Destination string
	TraceID     string
}

type GetEdgeParams

type GetEdgeParams struct {
	FromNodeID string
	ToNodeID   string
	Relation   string
}

type GetJobParams

type GetJobParams struct {
	JobID string
	Kinds []string
}

type GetJobTraceDeliveryParams

type GetJobTraceDeliveryParams struct {
	Destination string
	TraceID     string
}

type GetProjectionTaskParams

type GetProjectionTaskParams struct {
	SourceID string
	Revision string
}

type GetSourceRevisionParams

type GetSourceRevisionParams struct {
	SourceID string
	Revision string
}

type IDB

type IDB interface {
	Exec(ctx context.Context, sql string, arguments ...any) (pgconn.CommandTag, error)
	Query(ctx context.Context, sql string, arguments ...any) (pgx.Rows, error)
	QueryRow(ctx context.Context, sql string, arguments ...any) pgx.Row
	Begin(ctx context.Context) (pgx.Tx, error)
	BeginTx(ctx context.Context, options pgx.TxOptions) (pgx.Tx, error)
	Ping(ctx context.Context) error
}

IDB is the database capability a service receives. Its first three methods are sqlc's generated DBTX, so a service hands it straight to its generated queries; Begin and BeginTx open transactions and Ping reports readiness. *pgxpool.Pool, and therefore Pool, satisfies it.

type InsertIntentParams

type InsertIntentParams struct {
	IntentID  string
	SliceID   *string
	Intent    []byte
	CreatedAt pgtype.Timestamptz
}

type InsertJobMeasurementParams

type InsertJobMeasurementParams struct {
	JobID        string
	Metric       string
	Step         int64
	Value        float64
	RecordedAt   pgtype.Timestamptz
	EvidenceHash string
}

type InsertJobMetricDefinitionParams

type InsertJobMetricDefinitionParams struct {
	JobID       string
	Name        string
	Unit        string
	Description string
}

type InsertJobParams

type InsertJobParams struct {
	JobID                string
	Kind                 string
	Executor             string
	BudgetAccount        string
	Request              []byte
	RequestSha256        string
	TotalUnits           int64
	Managed              bool
	ExecutorTarget       string
	ExecutorImage        string
	ArtifactUri          string
	ReservationUsdMicros int64
	TimeoutSeconds       int32
}

type InsertRunningCronOccurrenceParams

type InsertRunningCronOccurrenceParams struct {
	OccurrenceID string
	TriggerName  string
	ScheduledAt  pgtype.Timestamptz
	LeaseOwner   *string
	LeaseToken   *string
	LeaseUntil   pgtype.Timestamptz
	ClaimedAt    pgtype.Timestamptz
}

type InsertSkippedCronOccurrenceParams

type InsertSkippedCronOccurrenceParams struct {
	OccurrenceID string
	TriggerName  string
	ScheduledAt  pgtype.Timestamptz
	SkippedAt    pgtype.Timestamptz
	SkipReason   *string
}

type InsertSliceEdgeParams

type InsertSliceEdgeParams struct {
	FromSliceID string
	ToSliceID   string
	Relation    string
}

type InsertSliceParams

type InsertSliceParams struct {
	SliceID    string
	Sequence   int64
	Title      string
	Recipe     []byte
	TouchSet   []byte
	Provenance []byte
	State      CsfSliceState
	CreatedAt  pgtype.Timestamptz
}

type LatestJobProgressRow

type LatestJobProgressRow struct {
	Kind           string
	Executor       string
	TotalUnits     int64
	CompletedUnits int64
	UpdatedAt      pgtype.Timestamptz
}

type ListChildNodesParams

type ListChildNodesParams struct {
	ParentNodeID *string
	ResultOffset int64
	ResultLimit  int32
}

type ListDocumentNodesParams

type ListDocumentNodesParams struct {
	ContentHash  *string
	ResultOffset int64
	ResultLimit  int32
}

type ListDocumentSourcesParams

type ListDocumentSourcesParams struct {
	ContentHash  string
	ResultOffset int64
	ResultLimit  int32
}

type ListDocumentsParams

type ListDocumentsParams struct {
	ResultOffset int64
	ResultLimit  int32
}

type ListExpiredCronOccurrencesParams

type ListExpiredCronOccurrencesParams struct {
	ExpiredAt pgtype.Timestamptz
	RowLimit  int32
}

type ListJobsParams

type ListJobsParams struct {
	Kinds    []string
	RowLimit int32
}

type ListNodeEdgesParams

type ListNodeEdgesParams struct {
	NodeID       string
	ResultOffset int64
	ResultLimit  int32
}

type ListSourceRevisionsParams

type ListSourceRevisionsParams struct {
	SourceID     string
	ResultOffset int64
	ResultLimit  int32
}

type ListSymbolNodesParams

type ListSymbolNodesParams struct {
	SymbolKey    string
	ResultOffset int64
	ResultLimit  int32
}

type LiveCronOccurrenceParams

type LiveCronOccurrenceParams struct {
	TriggerName          string
	At                   pgtype.Timestamptz
	ExcludedOccurrenceID string
}

type LockJobParams

type LockJobParams struct {
	JobID string
	Kinds []string
}

type MarkInterruptedJobSubmissionsParams

type MarkInterruptedJobSubmissionsParams struct {
	Reason       string
	Executor     string
	GraceSeconds int32
}

type NextHostJobParams

type NextHostJobParams struct {
	Executor string
	Kinds    []string
}

type NextUnarchivedJobParams

type NextUnarchivedJobParams struct {
	Executor string
	Kinds    []string
}

type NextUntracedJobParams

type NextUntracedJobParams struct {
	Executor string
	Kinds    []string
}

type NullCsfJobState

type NullCsfJobState struct {
	CsfJobState CsfJobState
	Valid       bool // Valid is true if CsfJobState is not NULL
}

func (*NullCsfJobState) Scan

func (ns *NullCsfJobState) Scan(value interface{}) error

Scan implements the Scanner interface.

func (NullCsfJobState) Value

func (ns NullCsfJobState) Value() (driver.Value, error)

Value implements the driver Valuer interface.

type NullCsfProjectionStatus

type NullCsfProjectionStatus struct {
	CsfProjectionStatus CsfProjectionStatus
	Valid               bool // Valid is true if CsfProjectionStatus is not NULL
}

func (*NullCsfProjectionStatus) Scan

func (ns *NullCsfProjectionStatus) Scan(value interface{}) error

Scan implements the Scanner interface.

func (NullCsfProjectionStatus) Value

func (ns NullCsfProjectionStatus) Value() (driver.Value, error)

Value implements the driver Valuer interface.

type NullCsfSliceState

type NullCsfSliceState struct {
	CsfSliceState CsfSliceState
	Valid         bool // Valid is true if CsfSliceState is not NULL
}

func (*NullCsfSliceState) Scan

func (ns *NullCsfSliceState) Scan(value interface{}) error

Scan implements the Scanner interface.

func (NullCsfSliceState) Value

func (ns NullCsfSliceState) Value() (driver.Value, error)

Value implements the driver Valuer interface.

type NullCsfTraceDeliveryState

type NullCsfTraceDeliveryState struct {
	CsfTraceDeliveryState CsfTraceDeliveryState
	Valid                 bool // Valid is true if CsfTraceDeliveryState is not NULL
}

func (*NullCsfTraceDeliveryState) Scan

func (ns *NullCsfTraceDeliveryState) Scan(value interface{}) error

Scan implements the Scanner interface.

func (NullCsfTraceDeliveryState) Value

Value implements the driver Valuer interface.

type PollableJobsParams

type PollableJobsParams struct {
	Executor string
	Kinds    []string
	RowLimit int32
}

type Pool

type Pool struct {
	*pgxpool.Pool
}

Pool is a PostgreSQL connection pool owned by the binary that opened it. It embeds the pgxpool.Pool, so it satisfies IDB; the binary alone calls Close, after every service it handed the pool to has stopped.

func OpenPool

func OpenPool(ctx context.Context, settings Settings) (*Pool, error)

OpenPool opens a pool for settings and verifies that the database answers before returning it. A pool that cannot reach the database is closed and not returned.

func (*Pool) OpenSQL

func (pool *Pool) OpenSQL() *sql.DB

OpenSQL returns a database/sql handle over the same connections, for the consumers built on database/sql, such as the shared migration runner candace/pkg/sqlmigrate. The caller closes the handle before the pool.

type Queries

type Queries struct {
	// contains filtered or unexported fields
}

func New

func New(db DBTX) *Queries

func (*Queries) AcquireExpiredCronOccurrence

func (q *Queries) AcquireExpiredCronOccurrence(ctx context.Context, arg AcquireExpiredCronOccurrenceParams) (CsfCronOccurrence, error)

A running occurrence whose lease expired is reclaimed as the next attempt.

func (*Queries) AdvanceCronTrigger

func (q *Queries) AdvanceCronTrigger(ctx context.Context, arg AdvanceCronTriggerParams) (int64, error)

The cursor only moves forward.

func (*Queries) BackfillProjectionTasks

func (q *Queries) BackfillProjectionTasks(ctx context.Context) (int64, error)

func (*Queries) ClaimJobSubmission

func (q *Queries) ClaimJobSubmission(ctx context.Context, arg ClaimJobSubmissionParams) (CsfJob, error)

Remote executors: a submission is claimed and committed before the executor is called, because a provider without an idempotency key cannot be retried blindly.

func (*Queries) ClaimJobTraceDelivery

func (q *Queries) ClaimJobTraceDelivery(ctx context.Context, arg ClaimJobTraceDeliveryParams) (string, error)

func (*Queries) ClaimProjectionTask

func (q *Queries) ClaimProjectionTask(ctx context.Context, arg ClaimProjectionTaskParams) (ClaimProjectionTaskRow, error)

Claim at most one due row. Expired final attempts become failed without granting another lease; no returned row is not proof that the queue is empty. Commit before calling the external projection backend.

func (*Queries) CompleteProjectionTask

func (q *Queries) CompleteProjectionTask(ctx context.Context, arg CompleteProjectionTaskParams) (CsfProjectionTask, error)

The lease must still be current and unexpired. A stale worker returns no row.

func (*Queries) CountJobStates

func (q *Queries) CountJobStates(ctx context.Context, kinds []string) ([]CountJobStatesRow, error)

func (*Queries) CountProjectionTasks

func (q *Queries) CountProjectionTasks(ctx context.Context) ([]CountProjectionTasksRow, error)

Include zero counts; the PostgreSQL enum owns the complete status set.

func (*Queries) CreateAgentConfiguration

func (q *Queries) CreateAgentConfiguration(ctx context.Context, arg CreateAgentConfigurationParams) (CsfAgentConfiguration, error)

An initial create is admitted exactly once. A duplicate reports no row so the service can return a revision conflict without overwriting state.

func (*Queries) CreateDocument

func (q *Queries) CreateDocument(ctx context.Context, arg CreateDocumentParams) (CsfDocument, error)

func (*Queries) CreateEdge

func (q *Queries) CreateEdge(ctx context.Context, arg CreateEdgeParams) (CsfEdge, error)

Supersession links versions of the same typed symbol. A refutation is an independent assertion and does not erase or deactivate either endpoint.

func (*Queries) CreateNode

func (q *Queries) CreateNode(ctx context.Context, arg CreateNodeParams) (CsfNode, error)

Parent creation precedes child creation; this API never changes a parent. Citation offsets are half-open byte intervals in the retained document.

func (*Queries) CreateSourceRevision

func (q *Queries) CreateSourceRevision(ctx context.Context, arg CreateSourceRevisionParams) (CsfSourceRevision, error)

Same immutable revision/content/metadata retains the first URI and retrieval time, even when a retry arrives through another locator at a later time.

func (*Queries) DisableCronTrigger

func (q *Queries) DisableCronTrigger(ctx context.Context, arg DisableCronTriggerParams) error

A trigger no longer declared keeps its history and stops firing.

func (*Queries) EnqueueProjectionTask

func (q *Queries) EnqueueProjectionTask(ctx context.Context, arg EnqueueProjectionTaskParams) (CsfProjectionTask, error)

Call in the source-registration transaction. Duplicate enqueue preserves the lease, retry schedule and terminal state; it never starts a second task.

func (*Queries) EnsureJobBudget

func (q *Queries) EnsureJobBudget(ctx context.Context, arg EnsureJobBudgetParams) (CsfJobBudget, error)

The job ledger: admitted requests to external executors. Every read that returns jobs is restricted to the kinds its caller decodes, so one ledger never hands another consumer's request to the wrong decoder. Budget accounts are immutable after creation: a repeated ensure with a different limit returns no row.

func (*Queries) FailProjectionTask

func (q *Queries) FailProjectionTask(ctx context.Context, arg FailProjectionTaskParams) (CsfProjectionTask, error)

Retry delay is min(max, base * 2^(attempts-1)); the exponent is bounded before evaluating POWER and both input delays are at most one day. No sleeping lease occupies a worker slot. A terminal task is not reset by duplicate ingestion.

func (*Queries) FinishCronOccurrence

func (q *Queries) FinishCronOccurrence(ctx context.Context, arg FinishCronOccurrenceParams) (int64, error)

func (*Queries) FinishJobSubmission

func (q *Queries) FinishJobSubmission(ctx context.Context, arg FinishJobSubmissionParams) (int64, error)

func (*Queries) FinishJobTraceDelivery

func (q *Queries) FinishJobTraceDelivery(ctx context.Context, arg FinishJobTraceDeliveryParams) (string, error)

func (*Queries) GetAgentConfiguration

func (q *Queries) GetAgentConfiguration(ctx context.Context, agentID string) (CsfAgentConfiguration, error)

func (*Queries) GetCronOccurrence

func (q *Queries) GetCronOccurrence(ctx context.Context, occurrenceID string) (CsfCronOccurrence, error)

func (*Queries) GetCronTrigger

func (q *Queries) GetCronTrigger(ctx context.Context, triggerName string) (CsfCronTrigger, error)

func (*Queries) GetDocument

func (q *Queries) GetDocument(ctx context.Context, contentHash string) (CsfDocument, error)

func (*Queries) GetEdge

func (q *Queries) GetEdge(ctx context.Context, arg GetEdgeParams) (CsfEdge, error)

func (*Queries) GetJob

func (q *Queries) GetJob(ctx context.Context, arg GetJobParams) (CsfJob, error)

func (*Queries) GetJobBudget

func (q *Queries) GetJobBudget(ctx context.Context, account string) (CsfJobBudget, error)

func (*Queries) GetJobTraceDelivery

func (q *Queries) GetJobTraceDelivery(ctx context.Context, arg GetJobTraceDeliveryParams) (CsfJobTraceDelivery, error)

func (*Queries) GetNode

func (q *Queries) GetNode(ctx context.Context, nodeID string) (CsfNode, error)

func (*Queries) GetProjectionTask

func (q *Queries) GetProjectionTask(ctx context.Context, arg GetProjectionTaskParams) (CsfProjectionTask, error)

func (*Queries) GetSourceRevision

func (q *Queries) GetSourceRevision(ctx context.Context, arg GetSourceRevisionParams) (CsfSourceRevision, error)

func (*Queries) InsertIntent

func (q *Queries) InsertIntent(ctx context.Context, arg InsertIntentParams) (CsfIntent, error)

func (*Queries) InsertJob

func (q *Queries) InsertJob(ctx context.Context, arg InsertJobParams) (CsfJob, error)

func (*Queries) InsertJobMeasurement

func (q *Queries) InsertJobMeasurement(ctx context.Context, arg InsertJobMeasurementParams) (int64, error)

A replayed identical value counts one row; a conflicting one counts none.

func (*Queries) InsertJobMetricDefinition

func (q *Queries) InsertJobMetricDefinition(ctx context.Context, arg InsertJobMetricDefinitionParams) (int64, error)

A replayed identical definition counts one row; a changed one counts none.

func (*Queries) InsertRunningCronOccurrence

func (q *Queries) InsertRunningCronOccurrence(ctx context.Context, arg InsertRunningCronOccurrenceParams) (CsfCronOccurrence, error)

func (*Queries) InsertSkippedCronOccurrence

func (q *Queries) InsertSkippedCronOccurrence(ctx context.Context, arg InsertSkippedCronOccurrenceParams) (CsfCronOccurrence, error)

func (*Queries) InsertSlice

func (q *Queries) InsertSlice(ctx context.Context, arg InsertSliceParams) (CsfSlice, error)

The slice graph: the dispatch service's persisted work queue. The service owns the graph in process and writes every change through; on start it reads the whole graph back. No read filters by caller: one dispatch service owns one harness process.

func (*Queries) InsertSliceEdge

func (q *Queries) InsertSliceEdge(ctx context.Context, arg InsertSliceEdgeParams) error

func (*Queries) LatestJobMeasurements

func (q *Queries) LatestJobMeasurements(ctx context.Context, jobID string) ([]CsfJobMeasurement, error)

func (*Queries) LatestJobProgress

func (q *Queries) LatestJobProgress(ctx context.Context, kinds []string) ([]LatestJobProgressRow, error)

func (*Queries) ListChildNodes

func (q *Queries) ListChildNodes(ctx context.Context, arg ListChildNodesParams) ([]CsfNode, error)

NULL selects roots, otherwise immediate children in the symbolic tree.

func (*Queries) ListCronTriggers

func (q *Queries) ListCronTriggers(ctx context.Context) ([]CsfCronTrigger, error)

Cron: the declared triggers and their occurrences; candace/services/cron owns the rules. Every write that fences on a lease is one conditional statement, so the store takes no row lock: a stale token or an expired lease changes no row, and the command tag reports it. Every time is a parameter the service read from its clock.

func (*Queries) ListDocumentNodes

func (q *Queries) ListDocumentNodes(ctx context.Context, arg ListDocumentNodesParams) ([]CsfNode, error)

func (*Queries) ListDocumentSources

func (q *Queries) ListDocumentSources(ctx context.Context, arg ListDocumentSourcesParams) ([]CsfSourceRevision, error)

func (*Queries) ListDocuments

func (q *Queries) ListDocuments(ctx context.Context, arg ListDocumentsParams) ([]CsfDocument, error)

func (*Queries) ListExpiredCronOccurrences

func (q *Queries) ListExpiredCronOccurrences(ctx context.Context, arg ListExpiredCronOccurrencesParams) ([]CsfCronOccurrence, error)

Abandoned leases of enabled triggers, oldest expiry first.

func (*Queries) ListIntents

func (q *Queries) ListIntents(ctx context.Context) ([]CsfIntent, error)

func (*Queries) ListJobMetricDefinitions

func (q *Queries) ListJobMetricDefinitions(ctx context.Context, jobID string) ([]CsfJobMetricDefinition, error)

func (*Queries) ListJobs

func (q *Queries) ListJobs(ctx context.Context, arg ListJobsParams) ([]CsfJob, error)

func (*Queries) ListNodeEdges

func (q *Queries) ListNodeEdges(ctx context.Context, arg ListNodeEdgesParams) ([]CsfEdge, error)

func (*Queries) ListRecentCronOccurrences

func (q *Queries) ListRecentCronOccurrences(ctx context.Context, rowLimit int32) ([]CsfCronOccurrence, error)

The newest occurrences, newest first; the caller reverses them.

func (*Queries) ListSliceEdges

func (q *Queries) ListSliceEdges(ctx context.Context) ([]CsfSliceEdge, error)

func (*Queries) ListSlices

func (q *Queries) ListSlices(ctx context.Context) ([]CsfSlice, error)

func (*Queries) ListSourceRevisions

func (q *Queries) ListSourceRevisions(ctx context.Context, arg ListSourceRevisionsParams) ([]CsfSourceRevision, error)

func (*Queries) ListSymbolNodes

func (q *Queries) ListSymbolNodes(ctx context.Context, arg ListSymbolNodesParams) ([]CsfNode, error)

All versions are returned. Neither timestamps nor model edges pick a winner.

func (*Queries) LiveCronOccurrence

func (q *Queries) LiveCronOccurrence(ctx context.Context, arg LiveCronOccurrenceParams) (string, error)

Another occurrence of the same trigger holding a live lease, for the overlap policy.

func (*Queries) LockJob

func (q *Queries) LockJob(ctx context.Context, arg LockJobParams) (CsfJob, error)

func (*Queries) LockJobBudget

func (q *Queries) LockJobBudget(ctx context.Context, account string) (CsfJobBudget, error)

The first statement of every admission transaction: it serializes admissions against one account, including concurrent retries.

func (*Queries) LockJobExecutor

func (q *Queries) LockJobExecutor(ctx context.Context, executor string) (bool, error)

Host executors: one transaction-scoped lock per executor serializes host reconciliation, even across overlapping deploys.

func (*Queries) MarkInterruptedJobSubmissions

func (q *Queries) MarkInterruptedJobSubmissions(ctx context.Context, arg MarkInterruptedJobSubmissionsParams) error

func (*Queries) NextHostJob

func (q *Queries) NextHostJob(ctx context.Context, arg NextHostJobParams) (CsfJob, error)

func (*Queries) NextUnarchivedJob

func (q *Queries) NextUnarchivedJob(ctx context.Context, arg NextUnarchivedJobParams) (CsfJob, error)

func (*Queries) NextUntracedJob

func (q *Queries) NextUntracedJob(ctx context.Context, arg NextUntracedJobParams) (CsfJob, error)

func (*Queries) PollableJobs

func (q *Queries) PollableJobs(ctx context.Context, arg PollableJobsParams) ([]CsfJob, error)

Terminal jobs stay pollable briefly so trailing log pages are drained.

func (*Queries) RecordJobProgress

func (q *Queries) RecordJobProgress(ctx context.Context, arg RecordJobProgressParams) error

func (*Queries) RenewCronOccurrenceLease

func (q *Queries) RenewCronOccurrenceLease(ctx context.Context, arg RenewCronOccurrenceLeaseParams) (int64, error)

func (*Queries) RequestJobCancellation

func (q *Queries) RequestJobCancellation(ctx context.Context, arg RequestJobCancellationParams) (CsfJob, error)

func (*Queries) ReserveJobBudget

func (q *Queries) ReserveJobBudget(ctx context.Context, arg ReserveJobBudgetParams) (int64, error)

func (*Queries) SearchSourceRevisions

func (q *Queries) SearchSourceRevisions(ctx context.Context, arg SearchSourceRevisionsParams) ([]CsfSourceRevision, error)

A bounded source search is a metadata lookup, not a full-text/vector search.

func (*Queries) SetJobInspectionError

func (q *Queries) SetJobInspectionError(ctx context.Context, arg SetJobInspectionErrorParams) error

func (*Queries) SetJobLogArchive

func (q *Queries) SetJobLogArchive(ctx context.Context, arg SetJobLogArchiveParams) error

func (*Queries) SetJobLogCursor

func (q *Queries) SetJobLogCursor(ctx context.Context, arg SetJobLogCursorParams) error

func (*Queries) SetJobTrace

func (q *Queries) SetJobTrace(ctx context.Context, arg SetJobTraceParams) error

func (*Queries) SkipExpiredCronOccurrence

func (q *Queries) SkipExpiredCronOccurrence(ctx context.Context, arg SkipExpiredCronOccurrenceParams) (CsfCronOccurrence, error)

func (*Queries) UpdateAgentConfiguration

func (q *Queries) UpdateAgentConfiguration(ctx context.Context, arg UpdateAgentConfigurationParams) (CsfAgentConfiguration, error)

A nonzero expected revision is a compare-and-swap update. A stale revision reports no row and never changes the retained secret references.

func (*Queries) UpdateJobExecution

func (q *Queries) UpdateJobExecution(ctx context.Context, arg UpdateJobExecutionParams) error

The update time moves only when the state does.

func (*Queries) UpdateSliceState

func (q *Queries) UpdateSliceState(ctx context.Context, arg UpdateSliceStateParams) (CsfSlice, error)

func (*Queries) UpsertCronTrigger

func (q *Queries) UpsertCronTrigger(ctx context.Context, arg UpsertCronTriggerParams) (CsfCronTrigger, error)

A declared trigger is inserted or brought back to its declaration, enabled.

func (*Queries) WithTx

func (q *Queries) WithTx(tx pgx.Tx) *Queries

type RecordJobProgressParams

type RecordJobProgressParams struct {
	CompletedUnits int64
	State          CsfJobState
	Reason         string
	JobID          string
}

type RenewCronOccurrenceLeaseParams

type RenewCronOccurrenceLeaseParams struct {
	LeaseUntil   pgtype.Timestamptz
	RenewedAt    pgtype.Timestamptz
	OccurrenceID string
	LeaseToken   *string
}

type RequestJobCancellationParams

type RequestJobCancellationParams struct {
	JobID string
	Kinds []string
}

type ReserveJobBudgetParams

type ReserveJobBudgetParams struct {
	Amount  int64
	Account string
}

type SearchSourceRevisionsParams

type SearchSourceRevisionsParams struct {
	Query        string
	ResultOffset int64
	ResultLimit  int32
}

type SetJobInspectionErrorParams

type SetJobInspectionErrorParams struct {
	JobID           string
	InspectionError string
}

type SetJobLogArchiveParams

type SetJobLogArchiveParams struct {
	JobID              string
	LogDocumentID      string
	LogProjectionError string
	LogIndexedAt       pgtype.Timestamptz
}

type SetJobLogCursorParams

type SetJobLogCursorParams struct {
	JobID     string
	LogCursor string
}

type SetJobTraceParams

type SetJobTraceParams struct {
	JobID            string
	TraceUrl         string
	TraceExportError string
}

type Settings

type Settings struct {
	// URL is a PostgreSQL connection string or URL as pgx parses it.
	URL string `json:"url"`
}

Settings selects the database a binary opens. The JSON shape is the CSF database configuration file's.

func SettingsFromEnvironment

func SettingsFromEnvironment(environment config.Environment, name string) (Settings, error)

SettingsFromEnvironment reads the database URL from the named variable of the config capability. The binary declares the name.

type SkipExpiredCronOccurrenceParams

type SkipExpiredCronOccurrenceParams struct {
	SkippedAt    pgtype.Timestamptz
	SkipReason   *string
	OccurrenceID string
}

type UpdateAgentConfigurationParams

type UpdateAgentConfigurationParams struct {
	LangfuseEndpointUrl            string
	LangfusePublicKeySecretRef     string
	LangfuseSecretKeySecretRef     string
	OpensearchEndpointUrl          string
	OpensearchIndex                string
	OpensearchEmbeddingModel       string
	OpensearchCredentialsSecretRef string
	AgentID                        string
	ExpectedRevision               int64
}

type UpdateJobExecutionParams

type UpdateJobExecutionParams struct {
	JobID            string
	State            CsfJobState
	ExternalID       string
	Reason           string
	LogStream        string
	CleanupConfirmed bool
}

type UpdateSliceStateParams

type UpdateSliceStateParams struct {
	SliceID        string
	State          CsfSliceState
	AssignmentID   string
	PullRequestUrl string
	Attempts       int32
	Checkpoint     string
	Error          string
	UpdatedAt      pgtype.Timestamptz
}

type UpsertCronTriggerParams

type UpsertCronTriggerParams struct {
	TriggerName         string
	ScheduleKind        string
	LocalHour           *int16
	LocalMinute         *int16
	Weekday             *int16
	MonthDay            *int16
	IntervalNanoseconds *int64
	RawExpression       *string
	Timezone            string
	IntervalAnchorAt    pgtype.Timestamptz
	NextRunAt           pgtype.Timestamptz
	CatchUpPolicy       string
	OverlapPolicy       string
	UpdatedAt           pgtype.Timestamptz
}

Directories

Path Synopsis
Package mocks is a generated GoMock package.
Package mocks is a generated GoMock package.

Jump to

Keyboard shortcuts

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