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
- Variables
- func ApplySchema(ctx context.Context, database *sql.DB) error
- func Time(value pgtype.Timestamptz) time.Time
- func Timestamp(value time.Time) pgtype.Timestamptz
- func Transact(ctx context.Context, database IDB, work func(tx pgx.Tx) error) error
- type AcquireExpiredCronOccurrenceParams
- type AdvanceCronTriggerParams
- type ClaimJobSubmissionParams
- type ClaimJobTraceDeliveryParams
- type ClaimProjectionTaskParams
- type ClaimProjectionTaskRow
- type CompleteProjectionTaskParams
- type CountJobStatesRow
- type CountProjectionTasksRow
- type CreateAgentConfigurationParams
- type CreateDocumentParams
- type CreateEdgeParams
- type CreateNodeParams
- type CreateSourceRevisionParams
- type CsfAgentConfiguration
- type CsfCronOccurrence
- type CsfCronTrigger
- type CsfDocument
- type CsfEdge
- type CsfIntent
- type CsfJob
- type CsfJobBudget
- type CsfJobMeasurement
- type CsfJobMetricDefinition
- type CsfJobState
- type CsfJobTraceDelivery
- type CsfNode
- type CsfProjectionStatus
- type CsfProjectionTask
- type CsfSlice
- type CsfSliceEdge
- type CsfSliceState
- type CsfSourceRevision
- type CsfTraceDeliveryState
- type DBTX
- type DisableCronTriggerParams
- type EnqueueProjectionTaskParams
- type EnsureJobBudgetParams
- type FailProjectionTaskParams
- type FinishCronOccurrenceParams
- type FinishJobSubmissionParams
- type FinishJobTraceDeliveryParams
- type GetEdgeParams
- type GetJobParams
- type GetJobTraceDeliveryParams
- type GetProjectionTaskParams
- type GetSourceRevisionParams
- type IDB
- type InsertIntentParams
- type InsertJobMeasurementParams
- type InsertJobMetricDefinitionParams
- type InsertJobParams
- type InsertRunningCronOccurrenceParams
- type InsertSkippedCronOccurrenceParams
- type InsertSliceEdgeParams
- type InsertSliceParams
- type LatestJobProgressRow
- type ListChildNodesParams
- type ListDocumentNodesParams
- type ListDocumentSourcesParams
- type ListDocumentsParams
- type ListExpiredCronOccurrencesParams
- type ListJobsParams
- type ListNodeEdgesParams
- type ListSourceRevisionsParams
- type ListSymbolNodesParams
- type LiveCronOccurrenceParams
- type LockJobParams
- type MarkInterruptedJobSubmissionsParams
- type NextHostJobParams
- type NextUnarchivedJobParams
- type NextUntracedJobParams
- type NullCsfJobState
- type NullCsfProjectionStatus
- type NullCsfSliceState
- type NullCsfTraceDeliveryState
- type PollableJobsParams
- type Pool
- type Queries
- func (q *Queries) AcquireExpiredCronOccurrence(ctx context.Context, arg AcquireExpiredCronOccurrenceParams) (CsfCronOccurrence, error)
- func (q *Queries) AdvanceCronTrigger(ctx context.Context, arg AdvanceCronTriggerParams) (int64, error)
- func (q *Queries) BackfillProjectionTasks(ctx context.Context) (int64, error)
- func (q *Queries) ClaimJobSubmission(ctx context.Context, arg ClaimJobSubmissionParams) (CsfJob, error)
- func (q *Queries) ClaimJobTraceDelivery(ctx context.Context, arg ClaimJobTraceDeliveryParams) (string, error)
- func (q *Queries) ClaimProjectionTask(ctx context.Context, arg ClaimProjectionTaskParams) (ClaimProjectionTaskRow, error)
- func (q *Queries) CompleteProjectionTask(ctx context.Context, arg CompleteProjectionTaskParams) (CsfProjectionTask, error)
- func (q *Queries) CountJobStates(ctx context.Context, kinds []string) ([]CountJobStatesRow, error)
- func (q *Queries) CountProjectionTasks(ctx context.Context) ([]CountProjectionTasksRow, error)
- func (q *Queries) CreateAgentConfiguration(ctx context.Context, arg CreateAgentConfigurationParams) (CsfAgentConfiguration, error)
- func (q *Queries) CreateDocument(ctx context.Context, arg CreateDocumentParams) (CsfDocument, error)
- func (q *Queries) CreateEdge(ctx context.Context, arg CreateEdgeParams) (CsfEdge, error)
- func (q *Queries) CreateNode(ctx context.Context, arg CreateNodeParams) (CsfNode, error)
- func (q *Queries) CreateSourceRevision(ctx context.Context, arg CreateSourceRevisionParams) (CsfSourceRevision, error)
- func (q *Queries) DisableCronTrigger(ctx context.Context, arg DisableCronTriggerParams) error
- func (q *Queries) EnqueueProjectionTask(ctx context.Context, arg EnqueueProjectionTaskParams) (CsfProjectionTask, error)
- func (q *Queries) EnsureJobBudget(ctx context.Context, arg EnsureJobBudgetParams) (CsfJobBudget, error)
- func (q *Queries) FailProjectionTask(ctx context.Context, arg FailProjectionTaskParams) (CsfProjectionTask, error)
- func (q *Queries) FinishCronOccurrence(ctx context.Context, arg FinishCronOccurrenceParams) (int64, error)
- func (q *Queries) FinishJobSubmission(ctx context.Context, arg FinishJobSubmissionParams) (int64, error)
- func (q *Queries) FinishJobTraceDelivery(ctx context.Context, arg FinishJobTraceDeliveryParams) (string, error)
- func (q *Queries) GetAgentConfiguration(ctx context.Context, agentID string) (CsfAgentConfiguration, error)
- func (q *Queries) GetCronOccurrence(ctx context.Context, occurrenceID string) (CsfCronOccurrence, error)
- func (q *Queries) GetCronTrigger(ctx context.Context, triggerName string) (CsfCronTrigger, error)
- func (q *Queries) GetDocument(ctx context.Context, contentHash string) (CsfDocument, error)
- func (q *Queries) GetEdge(ctx context.Context, arg GetEdgeParams) (CsfEdge, error)
- func (q *Queries) GetJob(ctx context.Context, arg GetJobParams) (CsfJob, error)
- func (q *Queries) GetJobBudget(ctx context.Context, account string) (CsfJobBudget, error)
- func (q *Queries) GetJobTraceDelivery(ctx context.Context, arg GetJobTraceDeliveryParams) (CsfJobTraceDelivery, error)
- func (q *Queries) GetNode(ctx context.Context, nodeID string) (CsfNode, error)
- func (q *Queries) GetProjectionTask(ctx context.Context, arg GetProjectionTaskParams) (CsfProjectionTask, error)
- func (q *Queries) GetSourceRevision(ctx context.Context, arg GetSourceRevisionParams) (CsfSourceRevision, error)
- func (q *Queries) InsertIntent(ctx context.Context, arg InsertIntentParams) (CsfIntent, error)
- func (q *Queries) InsertJob(ctx context.Context, arg InsertJobParams) (CsfJob, error)
- func (q *Queries) InsertJobMeasurement(ctx context.Context, arg InsertJobMeasurementParams) (int64, error)
- func (q *Queries) InsertJobMetricDefinition(ctx context.Context, arg InsertJobMetricDefinitionParams) (int64, error)
- func (q *Queries) InsertRunningCronOccurrence(ctx context.Context, arg InsertRunningCronOccurrenceParams) (CsfCronOccurrence, error)
- func (q *Queries) InsertSkippedCronOccurrence(ctx context.Context, arg InsertSkippedCronOccurrenceParams) (CsfCronOccurrence, error)
- func (q *Queries) InsertSlice(ctx context.Context, arg InsertSliceParams) (CsfSlice, error)
- func (q *Queries) InsertSliceEdge(ctx context.Context, arg InsertSliceEdgeParams) error
- func (q *Queries) LatestJobMeasurements(ctx context.Context, jobID string) ([]CsfJobMeasurement, error)
- func (q *Queries) LatestJobProgress(ctx context.Context, kinds []string) ([]LatestJobProgressRow, error)
- func (q *Queries) ListChildNodes(ctx context.Context, arg ListChildNodesParams) ([]CsfNode, error)
- func (q *Queries) ListCronTriggers(ctx context.Context) ([]CsfCronTrigger, error)
- func (q *Queries) ListDocumentNodes(ctx context.Context, arg ListDocumentNodesParams) ([]CsfNode, error)
- func (q *Queries) ListDocumentSources(ctx context.Context, arg ListDocumentSourcesParams) ([]CsfSourceRevision, error)
- func (q *Queries) ListDocuments(ctx context.Context, arg ListDocumentsParams) ([]CsfDocument, error)
- func (q *Queries) ListExpiredCronOccurrences(ctx context.Context, arg ListExpiredCronOccurrencesParams) ([]CsfCronOccurrence, error)
- func (q *Queries) ListIntents(ctx context.Context) ([]CsfIntent, error)
- func (q *Queries) ListJobMetricDefinitions(ctx context.Context, jobID string) ([]CsfJobMetricDefinition, error)
- func (q *Queries) ListJobs(ctx context.Context, arg ListJobsParams) ([]CsfJob, error)
- func (q *Queries) ListNodeEdges(ctx context.Context, arg ListNodeEdgesParams) ([]CsfEdge, error)
- func (q *Queries) ListRecentCronOccurrences(ctx context.Context, rowLimit int32) ([]CsfCronOccurrence, error)
- func (q *Queries) ListSliceEdges(ctx context.Context) ([]CsfSliceEdge, error)
- func (q *Queries) ListSlices(ctx context.Context) ([]CsfSlice, error)
- func (q *Queries) ListSourceRevisions(ctx context.Context, arg ListSourceRevisionsParams) ([]CsfSourceRevision, error)
- func (q *Queries) ListSymbolNodes(ctx context.Context, arg ListSymbolNodesParams) ([]CsfNode, error)
- func (q *Queries) LiveCronOccurrence(ctx context.Context, arg LiveCronOccurrenceParams) (string, error)
- func (q *Queries) LockJob(ctx context.Context, arg LockJobParams) (CsfJob, error)
- func (q *Queries) LockJobBudget(ctx context.Context, account string) (CsfJobBudget, error)
- func (q *Queries) LockJobExecutor(ctx context.Context, executor string) (bool, error)
- func (q *Queries) MarkInterruptedJobSubmissions(ctx context.Context, arg MarkInterruptedJobSubmissionsParams) error
- func (q *Queries) NextHostJob(ctx context.Context, arg NextHostJobParams) (CsfJob, error)
- func (q *Queries) NextUnarchivedJob(ctx context.Context, arg NextUnarchivedJobParams) (CsfJob, error)
- func (q *Queries) NextUntracedJob(ctx context.Context, arg NextUntracedJobParams) (CsfJob, error)
- func (q *Queries) PollableJobs(ctx context.Context, arg PollableJobsParams) ([]CsfJob, error)
- func (q *Queries) RecordJobProgress(ctx context.Context, arg RecordJobProgressParams) error
- func (q *Queries) RenewCronOccurrenceLease(ctx context.Context, arg RenewCronOccurrenceLeaseParams) (int64, error)
- func (q *Queries) RequestJobCancellation(ctx context.Context, arg RequestJobCancellationParams) (CsfJob, error)
- func (q *Queries) ReserveJobBudget(ctx context.Context, arg ReserveJobBudgetParams) (int64, error)
- func (q *Queries) SearchSourceRevisions(ctx context.Context, arg SearchSourceRevisionsParams) ([]CsfSourceRevision, error)
- func (q *Queries) SetJobInspectionError(ctx context.Context, arg SetJobInspectionErrorParams) error
- func (q *Queries) SetJobLogArchive(ctx context.Context, arg SetJobLogArchiveParams) error
- func (q *Queries) SetJobLogCursor(ctx context.Context, arg SetJobLogCursorParams) error
- func (q *Queries) SetJobTrace(ctx context.Context, arg SetJobTraceParams) error
- func (q *Queries) SkipExpiredCronOccurrence(ctx context.Context, arg SkipExpiredCronOccurrenceParams) (CsfCronOccurrence, error)
- func (q *Queries) UpdateAgentConfiguration(ctx context.Context, arg UpdateAgentConfigurationParams) (CsfAgentConfiguration, error)
- func (q *Queries) UpdateJobExecution(ctx context.Context, arg UpdateJobExecutionParams) error
- func (q *Queries) UpdateSliceState(ctx context.Context, arg UpdateSliceStateParams) (CsfSlice, error)
- func (q *Queries) UpsertCronTrigger(ctx context.Context, arg UpsertCronTriggerParams) (CsfCronTrigger, error)
- func (q *Queries) WithTx(tx pgx.Tx) *Queries
- type RecordJobProgressParams
- type RenewCronOccurrenceLeaseParams
- type RequestJobCancellationParams
- type ReserveJobBudgetParams
- type SearchSourceRevisionsParams
- type SetJobInspectionErrorParams
- type SetJobLogArchiveParams
- type SetJobLogCursorParams
- type SetJobTraceParams
- type Settings
- type SkipExpiredCronOccurrenceParams
- type UpdateAgentConfigurationParams
- type UpdateJobExecutionParams
- type UpdateSliceStateParams
- type UpsertCronTriggerParams
Constants ¶
const SchemaDirectory = "schema"
SchemaDirectory is the directory of Schema that holds the migrations.
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 ¶
var ErrMissingURL = errors.New("ipc/db/csfpg: a database URL is required")
ErrMissingURL means the settings name no database.
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 ¶
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.
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 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 CountJobStatesRow ¶
type CountJobStatesRow struct {
Kind string
Executor string
State CsfJobState
Jobs int64
}
type CountProjectionTasksRow ¶
type CountProjectionTasksRow struct {
Status CsfProjectionStatus
Count int64
}
type CreateDocumentParams ¶
type CreateEdgeParams ¶
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 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 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 CsfJobMetricDefinition ¶
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 CsfSliceEdge ¶
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 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 DisableCronTriggerParams ¶
type DisableCronTriggerParams struct {
UpdatedAt pgtype.Timestamptz
TriggerName string
}
type EnsureJobBudgetParams ¶
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 GetJobParams ¶
type GetProjectionTaskParams ¶
type GetSourceRevisionParams ¶
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 InsertJobParams ¶
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 InsertSliceParams ¶
type InsertSliceParams struct {
SliceID string
Sequence int64
Title string
Recipe []byte
TouchSet []byte
Provenance []byte
State CsfSliceState
CreatedAt pgtype.Timestamptz
}
type LatestJobProgressRow ¶
type ListChildNodesParams ¶
type ListDocumentNodesParams ¶
type ListDocumentsParams ¶
type ListExpiredCronOccurrencesParams ¶
type ListExpiredCronOccurrencesParams struct {
ExpiredAt pgtype.Timestamptz
RowLimit int32
}
type ListJobsParams ¶
type ListNodeEdgesParams ¶
type ListSymbolNodesParams ¶
type LiveCronOccurrenceParams ¶
type LiveCronOccurrenceParams struct {
TriggerName string
At pgtype.Timestamptz
ExcludedOccurrenceID string
}
type LockJobParams ¶
type NextHostJobParams ¶
type NextUnarchivedJobParams ¶
type NextUntracedJobParams ¶
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.
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.
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.
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.
type PollableJobsParams ¶
type 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.
type Queries ¶
type Queries struct {
// contains filtered or unexported fields
}
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 (*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 (*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 (*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 ¶
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 ¶
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 (*Queries) FinishJobSubmission ¶
func (*Queries) FinishJobTraceDelivery ¶
func (*Queries) GetAgentConfiguration ¶
func (*Queries) GetCronOccurrence ¶
func (*Queries) GetCronTrigger ¶
func (*Queries) GetDocument ¶
func (*Queries) GetJobBudget ¶
func (*Queries) GetJobTraceDelivery ¶
func (q *Queries) GetJobTraceDelivery(ctx context.Context, arg GetJobTraceDeliveryParams) (CsfJobTraceDelivery, 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 (*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 ¶
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 (*Queries) LatestJobProgress ¶
func (*Queries) ListChildNodes ¶
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 (*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 (*Queries) ListJobMetricDefinitions ¶
func (*Queries) ListNodeEdges ¶
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) 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) LockJobBudget ¶
The first statement of every admission transaction: it serializes admissions against one account, including concurrent retries.
func (*Queries) LockJobExecutor ¶
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 (*Queries) NextUnarchivedJob ¶
func (*Queries) NextUntracedJob ¶
func (*Queries) PollableJobs ¶
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 (*Queries) RequestJobCancellation ¶
func (*Queries) ReserveJobBudget ¶
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 (*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.
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 ReserveJobBudgetParams ¶
type SetJobLogArchiveParams ¶
type SetJobLogArchiveParams struct {
JobID string
LogDocumentID string
LogProjectionError string
LogIndexedAt pgtype.Timestamptz
}
type SetJobLogCursorParams ¶
type SetJobTraceParams ¶
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 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
}