Documentation
¶
Overview ¶
Package riverdatabasesql bundles a River driver for Go's built-in database/sql, making it interoperable with ORMs like Bun and GORM. It's generally still powered under the hood by Pgx because it's the only maintained, fully functional Postgres driver in the Go ecosystem, but it uses some lib/pq constructs internally by virtue of being implemented with Sqlc.
Example (MigrateDatabaseSQL) ¶
Example_migrateDatabaseSQL demonstrates the use of River's Go migration API through Go's built-in database/sql package.
package main
import (
"context"
"database/sql"
"fmt"
"strings"
_ "github.com/jackc/pgx/v5/stdlib"
"github.com/riverqueue/river/riverdbtest"
"github.com/riverqueue/river/riverdriver/riverdatabasesql"
"github.com/riverqueue/river/rivermigrate"
"github.com/riverqueue/river/rivershared/riversharedtest"
"github.com/riverqueue/river/rivershared/util/testutil"
"github.com/riverqueue/river/rivershared/util/urlutil"
)
func main() {
ctx := context.Background()
db, err := sql.Open("pgx", urlutil.DatabaseSQLCompatibleURL(riversharedtest.TestDatabaseURL()))
if err != nil {
panic(err)
}
defer db.Close()
driver := riverdatabasesql.New(db)
migrator, err := rivermigrate.New(driver, &rivermigrate.Config{
// Test schema with no migrations for purposes of this test.
Schema: riverdbtest.TestSchema(ctx, testutil.PanicTB(), driver, &riverdbtest.TestSchemaOpts{Lines: []string{}}),
})
if err != nil {
panic(err)
}
printVersions := func(res *rivermigrate.MigrateResult) {
for _, version := range res.Versions {
fmt.Printf("Migrated [%s] version %d\n", strings.ToUpper(string(res.Direction)), version.Version)
}
}
// Migrate to version 3. An actual call may want to omit all MigrateOpts,
// which will default to applying all available up migrations.
res, err := migrator.Migrate(ctx, rivermigrate.DirectionUp, &rivermigrate.MigrateOpts{
TargetVersion: 3,
})
if err != nil {
panic(err)
}
printVersions(res)
// Migrate down by three steps. Down migrating defaults to running only one
// step unless overridden by an option like MaxSteps or TargetVersion.
res, err = migrator.Migrate(ctx, rivermigrate.DirectionDown, &rivermigrate.MigrateOpts{
MaxSteps: 3,
})
if err != nil {
panic(err)
}
printVersions(res)
}
Output: Migrated [UP] version 1 Migrated [UP] version 2 Migrated [UP] version 3 Migrated [DOWN] version 3 Migrated [DOWN] version 2 Migrated [DOWN] version 1
Index ¶
- type Driver
- func (d *Driver) ArgPlaceholder() string
- func (d *Driver) DatabaseName() string
- func (d *Driver) GetExecutor() riverdriver.Executor
- func (d *Driver) GetListener(params *riverdriver.GetListenenerParams) riverdriver.Listener
- func (d *Driver) GetMigrationDefaultLines() []string
- func (d *Driver) GetMigrationFS(line string) fs.FS
- func (d *Driver) GetMigrationLines() []string
- func (d *Driver) GetMigrationTruncateTables(line string, version int) []string
- func (d *Driver) PoolIsSet() bool
- func (d *Driver) PoolSet(dbPool any) error
- func (d *Driver) SQLFragmentColumnContainsAll(column, namedArg string, values []string) (string, any, error)
- func (d *Driver) SQLFragmentColumnContainsAny(column, namedArg string, values []string) (string, any, error)
- func (d *Driver) SQLFragmentColumnIn(column string, values any) (string, any, error)
- func (d *Driver) SupportsListenNotify() bool
- func (d *Driver) SupportsListener() bool
- func (d *Driver) TimePrecision() time.Duration
- func (d *Driver) UnwrapExecutor(tx *sql.Tx) riverdriver.ExecutorTx
- func (d *Driver) UnwrapTx(execTx riverdriver.ExecutorTx) *sql.Tx
- type Executor
- func (e *Executor) Begin(ctx context.Context) (riverdriver.ExecutorTx, error)
- func (e *Executor) ColumnExists(ctx context.Context, params *riverdriver.ColumnExistsParams) (bool, error)
- func (e *Executor) Exec(ctx context.Context, sql string, args ...any) error
- func (e *Executor) IndexDropIfExists(ctx context.Context, params *riverdriver.IndexDropIfExistsParams) error
- func (e *Executor) IndexExists(ctx context.Context, params *riverdriver.IndexExistsParams) (bool, error)
- func (e *Executor) IndexReindex(ctx context.Context, params *riverdriver.IndexReindexParams) error
- func (e *Executor) IndexReindexArtifacts(ctx context.Context, params *riverdriver.IndexReindexArtifactsParams) ([]string, error)
- func (e *Executor) IndexesExist(ctx context.Context, params *riverdriver.IndexesExistParams) (map[string]bool, error)
- func (e *Executor) JobCancel(ctx context.Context, params *riverdriver.JobCancelParams) (*rivertype.JobRow, error)
- func (e *Executor) JobCountByAllStates(ctx context.Context, params *riverdriver.JobCountByAllStatesParams) (map[rivertype.JobState]int, error)
- func (e *Executor) JobCountByQueueAndState(ctx context.Context, params *riverdriver.JobCountByQueueAndStateParams) ([]*riverdriver.JobCountByQueueAndStateResult, error)
- func (e *Executor) JobCountByState(ctx context.Context, params *riverdriver.JobCountByStateParams) (int, error)
- func (e *Executor) JobDelete(ctx context.Context, params *riverdriver.JobDeleteParams) (*rivertype.JobRow, error)
- func (e *Executor) JobDeleteBefore(ctx context.Context, params *riverdriver.JobDeleteBeforeParams) (int, error)
- func (e *Executor) JobDeleteMany(ctx context.Context, params *riverdriver.JobDeleteManyParams) ([]*rivertype.JobRow, error)
- func (e *Executor) JobGetAvailable(ctx context.Context, params *riverdriver.JobGetAvailableParams) ([]*rivertype.JobRow, error)
- func (e *Executor) JobGetByID(ctx context.Context, params *riverdriver.JobGetByIDParams) (*rivertype.JobRow, error)
- func (e *Executor) JobGetByIDMany(ctx context.Context, params *riverdriver.JobGetByIDManyParams) ([]*rivertype.JobRow, error)
- func (e *Executor) JobGetByKindMany(ctx context.Context, params *riverdriver.JobGetByKindManyParams) ([]*rivertype.JobRow, error)
- func (e *Executor) JobGetStuck(ctx context.Context, params *riverdriver.JobGetStuckParams) ([]*rivertype.JobRow, error)
- func (e *Executor) JobInsertFastMany(ctx context.Context, params *riverdriver.JobInsertFastManyParams) ([]*riverdriver.JobInsertFastResult, error)
- func (e *Executor) JobInsertFastManyNoReturning(ctx context.Context, params *riverdriver.JobInsertFastManyParams) (int, error)
- func (e *Executor) JobInsertFull(ctx context.Context, params *riverdriver.JobInsertFullParams) (*rivertype.JobRow, error)
- func (e *Executor) JobInsertFullMany(ctx context.Context, params *riverdriver.JobInsertFullManyParams) ([]*rivertype.JobRow, error)
- func (e *Executor) JobKindList(ctx context.Context, params *riverdriver.JobKindListParams) ([]string, error)
- func (e *Executor) JobList(ctx context.Context, params *riverdriver.JobListParams) ([]*rivertype.JobRow, error)
- func (e *Executor) JobRescueMany(ctx context.Context, params *riverdriver.JobRescueManyParams) (*struct{}, error)
- func (e *Executor) JobRetry(ctx context.Context, params *riverdriver.JobRetryParams) (*rivertype.JobRow, error)
- func (e *Executor) JobSchedule(ctx context.Context, params *riverdriver.JobScheduleParams) ([]*riverdriver.JobScheduleResult, error)
- func (e *Executor) JobSetStateIfRunningMany(ctx context.Context, params *riverdriver.JobSetStateIfRunningManyParams) ([]*rivertype.JobRow, error)
- func (e *Executor) JobUpdate(ctx context.Context, params *riverdriver.JobUpdateParams) (*rivertype.JobRow, error)
- func (e *Executor) JobUpdateFull(ctx context.Context, params *riverdriver.JobUpdateFullParams) (*rivertype.JobRow, error)
- func (e *Executor) LeaderAttemptElect(ctx context.Context, params *riverdriver.LeaderElectParams) (*riverdriver.Leader, error)
- func (e *Executor) LeaderAttemptReelect(ctx context.Context, params *riverdriver.LeaderReelectParams) (*riverdriver.Leader, error)
- func (e *Executor) LeaderDeleteExpired(ctx context.Context, params *riverdriver.LeaderDeleteExpiredParams) (int, error)
- func (e *Executor) LeaderGetElectedLeader(ctx context.Context, params *riverdriver.LeaderGetElectedLeaderParams) (*riverdriver.Leader, error)
- func (e *Executor) LeaderInsert(ctx context.Context, params *riverdriver.LeaderInsertParams) (*riverdriver.Leader, error)
- func (e *Executor) LeaderResign(ctx context.Context, params *riverdriver.LeaderResignParams) (bool, error)
- func (e *Executor) MigrationDeleteAssumingMainMany(ctx context.Context, params *riverdriver.MigrationDeleteAssumingMainManyParams) ([]*riverdriver.Migration, error)
- func (e *Executor) MigrationDeleteByLineAndVersionMany(ctx context.Context, ...) ([]*riverdriver.Migration, error)
- func (e *Executor) MigrationGetAllAssumingMain(ctx context.Context, params *riverdriver.MigrationGetAllAssumingMainParams) ([]*riverdriver.Migration, error)
- func (e *Executor) MigrationGetByLine(ctx context.Context, params *riverdriver.MigrationGetByLineParams) ([]*riverdriver.Migration, error)
- func (e *Executor) MigrationInsertMany(ctx context.Context, params *riverdriver.MigrationInsertManyParams) ([]*riverdriver.Migration, error)
- func (e *Executor) MigrationInsertManyAssumingMain(ctx context.Context, params *riverdriver.MigrationInsertManyAssumingMainParams) ([]*riverdriver.Migration, error)
- func (e *Executor) NotificationDeleteBefore(ctx context.Context, params *riverdriver.NotificationDeleteBeforeParams) (int, error)
- func (e *Executor) NotifyMany(ctx context.Context, params *riverdriver.NotifyManyParams) error
- func (e *Executor) PGAdvisoryXactLock(ctx context.Context, key int64) (*struct{}, error)
- func (e *Executor) QueryRow(ctx context.Context, sql string, args ...any) riverdriver.Row
- func (e *Executor) QueueCreateOrSetUpdatedAt(ctx context.Context, params *riverdriver.QueueCreateOrSetUpdatedAtParams) (*rivertype.Queue, error)
- func (e *Executor) QueueDeleteExpired(ctx context.Context, params *riverdriver.QueueDeleteExpiredParams) ([]string, error)
- func (e *Executor) QueueGet(ctx context.Context, params *riverdriver.QueueGetParams) (*rivertype.Queue, error)
- func (e *Executor) QueueList(ctx context.Context, params *riverdriver.QueueListParams) ([]*rivertype.Queue, error)
- func (e *Executor) QueueNameList(ctx context.Context, params *riverdriver.QueueNameListParams) ([]string, error)
- func (e *Executor) QueuePause(ctx context.Context, params *riverdriver.QueuePauseParams) error
- func (e *Executor) QueueResume(ctx context.Context, params *riverdriver.QueueResumeParams) error
- func (e *Executor) QueueUpdate(ctx context.Context, params *riverdriver.QueueUpdateParams) (*rivertype.Queue, error)
- func (e *Executor) SchemaCreate(ctx context.Context, params *riverdriver.SchemaCreateParams) error
- func (e *Executor) SchemaDrop(ctx context.Context, params *riverdriver.SchemaDropParams) error
- func (e *Executor) SchemaGetExpired(ctx context.Context, params *riverdriver.SchemaGetExpiredParams) ([]string, error)
- func (e *Executor) TableExists(ctx context.Context, params *riverdriver.TableExistsParams) (bool, error)
- func (e *Executor) TableTruncate(ctx context.Context, params *riverdriver.TableTruncateParams) error
- type ExecutorSubTx
- type ExecutorTx
Examples ¶
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type Driver ¶
type Driver struct {
// contains filtered or unexported fields
}
Driver is an implementation of riverdriver.Driver for database/sql.
func New ¶
New returns a new database/sql River driver for use with River.
It takes an sql.DB to use for use with River. The pool should already be configured to use the schema specified in the client's Schema field. The pool must not be closed while associated River objects are running.
func NewWithPgxListener ¶ added in v0.46.0
NewWithPgxListener returns a new database/sql River driver with a Pgx-backed listener. The database/sql pool continues to be used for all database operations other than listening for notifications. The Pgx pool is used only to acquire dedicated connections for Postgres LISTEN commands. It panics if listenerPool is nil; use New for a poll-only driver.
Both pools are owned by the caller, must connect to the same database, and must resolve the same schema. When the River client has no explicit schema, both pools' search paths must produce the same current schema. Neither pool may be closed while associated River objects are running.
Listener connections are hijacked from listenerPool and never returned. A pool dedicated to one River client can generally set MinConns to zero and MaxConns to one. Each concurrently running client still needs its own listener connection. Because hijacked connections no longer count against the pool's maximum, sharing a listener pool between clients may cause total live connections to exceed that maximum. Closing listenerPool does not close hijacked connections; stopping the associated River clients does.
Applications using PgBouncer must configure the listener pool to use session pooling or connect it directly to Postgres.
func (*Driver) ArgPlaceholder ¶ added in v0.23.0
func (*Driver) DatabaseName ¶ added in v0.23.0
func (*Driver) GetExecutor ¶
func (d *Driver) GetExecutor() riverdriver.Executor
func (*Driver) GetListener ¶ added in v0.0.23
func (d *Driver) GetListener(params *riverdriver.GetListenenerParams) riverdriver.Listener
func (*Driver) GetMigrationDefaultLines ¶ added in v0.21.0
func (*Driver) GetMigrationLines ¶ added in v0.10.0
func (*Driver) GetMigrationTruncateTables ¶ added in v0.21.0
func (*Driver) SQLFragmentColumnContainsAll ¶ added in v0.43.0
func (*Driver) SQLFragmentColumnContainsAny ¶ added in v0.43.0
func (*Driver) SQLFragmentColumnIn ¶ added in v0.23.0
func (*Driver) SupportsListenNotify ¶ added in v0.23.0
func (*Driver) SupportsListener ¶ added in v0.10.0
func (*Driver) TimePrecision ¶ added in v0.23.0
func (*Driver) UnwrapExecutor ¶
func (d *Driver) UnwrapExecutor(tx *sql.Tx) riverdriver.ExecutorTx
func (*Driver) UnwrapTx ¶
func (d *Driver) UnwrapTx(execTx riverdriver.ExecutorTx) *sql.Tx
type Executor ¶
type Executor struct {
// contains filtered or unexported fields
}
func (*Executor) Begin ¶
func (e *Executor) Begin(ctx context.Context) (riverdriver.ExecutorTx, error)
func (*Executor) ColumnExists ¶ added in v0.10.0
func (e *Executor) ColumnExists(ctx context.Context, params *riverdriver.ColumnExistsParams) (bool, error)
func (*Executor) IndexDropIfExists ¶ added in v0.23.0
func (e *Executor) IndexDropIfExists(ctx context.Context, params *riverdriver.IndexDropIfExistsParams) error
func (*Executor) IndexExists ¶ added in v0.23.0
func (e *Executor) IndexExists(ctx context.Context, params *riverdriver.IndexExistsParams) (bool, error)
func (*Executor) IndexReindex ¶ added in v0.23.0
func (e *Executor) IndexReindex(ctx context.Context, params *riverdriver.IndexReindexParams) error
func (*Executor) IndexReindexArtifacts ¶ added in v0.40.0
func (e *Executor) IndexReindexArtifacts(ctx context.Context, params *riverdriver.IndexReindexArtifactsParams) ([]string, error)
func (*Executor) IndexesExist ¶ added in v0.24.0
func (e *Executor) IndexesExist(ctx context.Context, params *riverdriver.IndexesExistParams) (map[string]bool, error)
func (*Executor) JobCancel ¶ added in v0.0.23
func (e *Executor) JobCancel(ctx context.Context, params *riverdriver.JobCancelParams) (*rivertype.JobRow, error)
func (*Executor) JobCountByAllStates ¶ added in v0.24.0
func (e *Executor) JobCountByAllStates(ctx context.Context, params *riverdriver.JobCountByAllStatesParams) (map[rivertype.JobState]int, error)
func (*Executor) JobCountByQueueAndState ¶ added in v0.24.0
func (e *Executor) JobCountByQueueAndState(ctx context.Context, params *riverdriver.JobCountByQueueAndStateParams) ([]*riverdriver.JobCountByQueueAndStateResult, error)
func (*Executor) JobCountByState ¶ added in v0.1.0
func (e *Executor) JobCountByState(ctx context.Context, params *riverdriver.JobCountByStateParams) (int, error)
func (*Executor) JobDelete ¶ added in v0.7.0
func (e *Executor) JobDelete(ctx context.Context, params *riverdriver.JobDeleteParams) (*rivertype.JobRow, error)
func (*Executor) JobDeleteBefore ¶ added in v0.0.23
func (e *Executor) JobDeleteBefore(ctx context.Context, params *riverdriver.JobDeleteBeforeParams) (int, error)
func (*Executor) JobDeleteMany ¶ added in v0.24.0
func (e *Executor) JobDeleteMany(ctx context.Context, params *riverdriver.JobDeleteManyParams) ([]*rivertype.JobRow, error)
func (*Executor) JobGetAvailable ¶ added in v0.0.23
func (e *Executor) JobGetAvailable(ctx context.Context, params *riverdriver.JobGetAvailableParams) ([]*rivertype.JobRow, error)
func (*Executor) JobGetByID ¶ added in v0.0.23
func (e *Executor) JobGetByID(ctx context.Context, params *riverdriver.JobGetByIDParams) (*rivertype.JobRow, error)
func (*Executor) JobGetByIDMany ¶ added in v0.0.23
func (e *Executor) JobGetByIDMany(ctx context.Context, params *riverdriver.JobGetByIDManyParams) ([]*rivertype.JobRow, error)
func (*Executor) JobGetByKindMany ¶ added in v0.0.23
func (e *Executor) JobGetByKindMany(ctx context.Context, params *riverdriver.JobGetByKindManyParams) ([]*rivertype.JobRow, error)
func (*Executor) JobGetStuck ¶ added in v0.0.23
func (e *Executor) JobGetStuck(ctx context.Context, params *riverdriver.JobGetStuckParams) ([]*rivertype.JobRow, error)
func (*Executor) JobInsertFastMany ¶ added in v0.0.23
func (e *Executor) JobInsertFastMany(ctx context.Context, params *riverdriver.JobInsertFastManyParams) ([]*riverdriver.JobInsertFastResult, error)
func (*Executor) JobInsertFastManyNoReturning ¶ added in v0.12.0
func (e *Executor) JobInsertFastManyNoReturning(ctx context.Context, params *riverdriver.JobInsertFastManyParams) (int, error)
func (*Executor) JobInsertFull ¶ added in v0.0.23
func (e *Executor) JobInsertFull(ctx context.Context, params *riverdriver.JobInsertFullParams) (*rivertype.JobRow, error)
func (*Executor) JobInsertFullMany ¶ added in v0.23.0
func (e *Executor) JobInsertFullMany(ctx context.Context, params *riverdriver.JobInsertFullManyParams) ([]*rivertype.JobRow, error)
func (*Executor) JobKindList ¶ added in v0.24.0
func (e *Executor) JobKindList(ctx context.Context, params *riverdriver.JobKindListParams) ([]string, error)
func (*Executor) JobList ¶ added in v0.0.23
func (e *Executor) JobList(ctx context.Context, params *riverdriver.JobListParams) ([]*rivertype.JobRow, error)
func (*Executor) JobRescueMany ¶ added in v0.0.23
func (e *Executor) JobRescueMany(ctx context.Context, params *riverdriver.JobRescueManyParams) (*struct{}, error)
func (*Executor) JobRetry ¶ added in v0.0.23
func (e *Executor) JobRetry(ctx context.Context, params *riverdriver.JobRetryParams) (*rivertype.JobRow, error)
func (*Executor) JobSchedule ¶ added in v0.0.23
func (e *Executor) JobSchedule(ctx context.Context, params *riverdriver.JobScheduleParams) ([]*riverdriver.JobScheduleResult, error)
func (*Executor) JobSetStateIfRunningMany ¶ added in v0.12.1
func (e *Executor) JobSetStateIfRunningMany(ctx context.Context, params *riverdriver.JobSetStateIfRunningManyParams) ([]*rivertype.JobRow, error)
func (*Executor) JobUpdate ¶ added in v0.0.23
func (e *Executor) JobUpdate(ctx context.Context, params *riverdriver.JobUpdateParams) (*rivertype.JobRow, error)
func (*Executor) JobUpdateFull ¶ added in v0.29.0
func (e *Executor) JobUpdateFull(ctx context.Context, params *riverdriver.JobUpdateFullParams) (*rivertype.JobRow, error)
func (*Executor) LeaderAttemptElect ¶ added in v0.0.23
func (e *Executor) LeaderAttemptElect(ctx context.Context, params *riverdriver.LeaderElectParams) (*riverdriver.Leader, error)
func (*Executor) LeaderAttemptReelect ¶ added in v0.0.23
func (e *Executor) LeaderAttemptReelect(ctx context.Context, params *riverdriver.LeaderReelectParams) (*riverdriver.Leader, error)
func (*Executor) LeaderDeleteExpired ¶ added in v0.0.23
func (e *Executor) LeaderDeleteExpired(ctx context.Context, params *riverdriver.LeaderDeleteExpiredParams) (int, error)
func (*Executor) LeaderGetElectedLeader ¶ added in v0.0.23
func (e *Executor) LeaderGetElectedLeader(ctx context.Context, params *riverdriver.LeaderGetElectedLeaderParams) (*riverdriver.Leader, error)
func (*Executor) LeaderInsert ¶ added in v0.0.23
func (e *Executor) LeaderInsert(ctx context.Context, params *riverdriver.LeaderInsertParams) (*riverdriver.Leader, error)
func (*Executor) LeaderResign ¶ added in v0.0.23
func (e *Executor) LeaderResign(ctx context.Context, params *riverdriver.LeaderResignParams) (bool, error)
func (*Executor) MigrationDeleteAssumingMainMany ¶ added in v0.10.0
func (e *Executor) MigrationDeleteAssumingMainMany(ctx context.Context, params *riverdriver.MigrationDeleteAssumingMainManyParams) ([]*riverdriver.Migration, error)
func (*Executor) MigrationDeleteByLineAndVersionMany ¶ added in v0.10.0
func (e *Executor) MigrationDeleteByLineAndVersionMany(ctx context.Context, params *riverdriver.MigrationDeleteByLineAndVersionManyParams) ([]*riverdriver.Migration, error)
func (*Executor) MigrationGetAllAssumingMain ¶ added in v0.10.0
func (e *Executor) MigrationGetAllAssumingMain(ctx context.Context, params *riverdriver.MigrationGetAllAssumingMainParams) ([]*riverdriver.Migration, error)
func (*Executor) MigrationGetByLine ¶ added in v0.10.0
func (e *Executor) MigrationGetByLine(ctx context.Context, params *riverdriver.MigrationGetByLineParams) ([]*riverdriver.Migration, error)
func (*Executor) MigrationInsertMany ¶
func (e *Executor) MigrationInsertMany(ctx context.Context, params *riverdriver.MigrationInsertManyParams) ([]*riverdriver.Migration, error)
func (*Executor) MigrationInsertManyAssumingMain ¶ added in v0.10.0
func (e *Executor) MigrationInsertManyAssumingMain(ctx context.Context, params *riverdriver.MigrationInsertManyAssumingMainParams) ([]*riverdriver.Migration, error)
func (*Executor) NotificationDeleteBefore ¶ added in v0.40.0
func (e *Executor) NotificationDeleteBefore(ctx context.Context, params *riverdriver.NotificationDeleteBeforeParams) (int, error)
func (*Executor) NotifyMany ¶ added in v0.5.0
func (e *Executor) NotifyMany(ctx context.Context, params *riverdriver.NotifyManyParams) error
func (*Executor) PGAdvisoryXactLock ¶ added in v0.0.23
func (*Executor) QueueCreateOrSetUpdatedAt ¶ added in v0.5.0
func (e *Executor) QueueCreateOrSetUpdatedAt(ctx context.Context, params *riverdriver.QueueCreateOrSetUpdatedAtParams) (*rivertype.Queue, error)
func (*Executor) QueueDeleteExpired ¶ added in v0.5.0
func (e *Executor) QueueDeleteExpired(ctx context.Context, params *riverdriver.QueueDeleteExpiredParams) ([]string, error)
func (*Executor) QueueGet ¶ added in v0.5.0
func (e *Executor) QueueGet(ctx context.Context, params *riverdriver.QueueGetParams) (*rivertype.Queue, error)
func (*Executor) QueueList ¶ added in v0.5.0
func (e *Executor) QueueList(ctx context.Context, params *riverdriver.QueueListParams) ([]*rivertype.Queue, error)
func (*Executor) QueueNameList ¶ added in v0.24.0
func (e *Executor) QueueNameList(ctx context.Context, params *riverdriver.QueueNameListParams) ([]string, error)
func (*Executor) QueuePause ¶ added in v0.5.0
func (e *Executor) QueuePause(ctx context.Context, params *riverdriver.QueuePauseParams) error
func (*Executor) QueueResume ¶ added in v0.5.0
func (e *Executor) QueueResume(ctx context.Context, params *riverdriver.QueueResumeParams) error
func (*Executor) QueueUpdate ¶ added in v0.20.0
func (e *Executor) QueueUpdate(ctx context.Context, params *riverdriver.QueueUpdateParams) (*rivertype.Queue, error)
func (*Executor) SchemaCreate ¶ added in v0.23.0
func (e *Executor) SchemaCreate(ctx context.Context, params *riverdriver.SchemaCreateParams) error
func (*Executor) SchemaDrop ¶ added in v0.23.0
func (e *Executor) SchemaDrop(ctx context.Context, params *riverdriver.SchemaDropParams) error
func (*Executor) SchemaGetExpired ¶ added in v0.21.0
func (e *Executor) SchemaGetExpired(ctx context.Context, params *riverdriver.SchemaGetExpiredParams) ([]string, error)
func (*Executor) TableExists ¶
func (e *Executor) TableExists(ctx context.Context, params *riverdriver.TableExistsParams) (bool, error)
func (*Executor) TableTruncate ¶ added in v0.23.0
func (e *Executor) TableTruncate(ctx context.Context, params *riverdriver.TableTruncateParams) error
type ExecutorSubTx ¶ added in v0.10.0
type ExecutorSubTx struct {
Executor
// contains filtered or unexported fields
}
func (*ExecutorSubTx) Begin ¶ added in v0.10.0
func (t *ExecutorSubTx) Begin(ctx context.Context) (riverdriver.ExecutorTx, error)
type ExecutorTx ¶
type ExecutorTx struct {
Executor
// contains filtered or unexported fields
}
func (*ExecutorTx) Begin ¶ added in v0.10.0
func (t *ExecutorTx) Begin(ctx context.Context) (riverdriver.ExecutorTx, error)