MCPcopy Create free account

hub / github.com/dataddo/pgq / functions

Functions122 in github.com/dataddo/pgq

↓ 73 callersFunctionNoError
NoError fails the test if err is not nil.
internal/require/require.go:11
↓ 34 callersMethodWriteString
WriteString appends the provided string to the query and extracts any parameters from it.
internal/query/query_builder.go:24
↓ 31 callersMethodExecContext
(ctx context.Context, query string, args ...interface{})
consumer.go:712
↓ 22 callersFunctionEqual
Equal fails the test if expected is not equal to actual.
internal/require/require.go:35
↓ 20 callersMethodString
()
internal/query/query_builder.go:42
↓ 18 callersFunctionQuoteIdentifier
QuoteIdentifier quotes an "identifier" (e.g. a table or a column name) to be used as part of an SQL statement. For example: tblname := "my_table"
internal/pg/pg.go:33
↓ 13 callersFunctionopenDB
(t *testing.T)
integtest/consumer_test.go:324
↓ 12 callersMethodRun
Run consumes messages until the context isn't cancelled.
consumer.go:318
↓ 10 callersFunctionGenerateDropTableQuery
GenerateDropTableQuery returns a postgres query for dropping the queue table
x/schema/queue.go:34
↓ 9 callersMethodPublish
(ctx context.Context, queue string, msg ...*MessageOutgoing)
publisher.go:25
↓ 8 callersFunctionNewPublisher
NewPublisher initializes the publisher with given *sql.DB client.
publisher.go:54
↓ 7 callersFunctionNewConsumer
NewConsumer creates Consumer with proper settings
consumer.go:264
↓ 7 callersMethodNext
Next returns next parameter.
internal/pg/pg.go:15
↓ 6 callersMethodError
()
errors.go:22
↓ 6 callersMethodFatal
()
errors.go:23
↓ 6 callersFunctiongenerateDropTableQuery
(queueName string)
validator_test.go:263
↓ 6 callersFunctiongenerateRandomString
(length int)
validator_test.go:183
↓ 6 callersFunctionopenDB
TODO: This was recovered from the consumer_test.go file. We can make a common testing package and add all these common functionalities will be include
validator_test.go:155
↓ 5 callersFunctionErrorIs
ErrorIs fails the test if err is nil or does not match target.
internal/require/require.go:27
↓ 5 callersFunctionGenerateCreateTableQuery
GenerateCreateTableQuery returns the query for creating the queue table
x/schema/queue.go:11
↓ 5 callersFunctionWithInvalidMessageCallback
WithInvalidMessageCallback sets callback for invalid messages.
consumer.go:196
↓ 5 callersFunctionWithLockDuration
WithLockDuration sets the maximal duration for how long the message remains locked for other consumers.
consumer.go:151
↓ 5 callersFunctionWithLogger
WithLogger sets logger. Default is no logging.
consumer.go:223
↓ 5 callersFunctionWithMaxParallelMessages
WithMaxParallelMessages sets how many jobs can single consumer process simultaneously.
consumer.go:175
↓ 5 callersFunctionWithMetrics
WithMetrics sets metrics meter. Default is noop.Meter{}.
consumer.go:182
↓ 5 callersFunctionWithPollingInterval
WithPollingInterval sets how frequently consumer checks the queue for new messages.
consumer.go:159
↓ 5 callersFunctionacquireMaxFromSemaphore
acquireMaxFromSemaphore acquires maximum possible weight. It blocks until resources are available or ctx is done. On success, returns acquired weight.
consumer.go:794
↓ 4 callersFunctionError
Error fails the test if err is nil.
internal/require/require.go:19
↓ 4 callersMethodSetDeadline
SetDeadline sets the message deadline. If the deadline is after the message deadline, it returns ErrInvalidDeadline. The timeout also affects the queu
message.go:129
↓ 4 callersFunctionValidateFields
ValidateFields checks if required fields exist
validator.go:55
↓ 4 callersFunctionValidateIndexes
ValidateIndexes checks if required indexes exist
validator.go:91
↓ 4 callersFunctionWithHistoryLimit
WithHistoryLimit sets how long in history you want to search for unprocessed messages (default is no limit). If not set, it will look for message in t
consumer.go:206
↓ 3 callersFunctionStaticMetaInjector
StaticMetaInjector returns a Metadata injector that injects given Metadata.
publisher.go:46
↓ 3 callersFunctionWithMetaInjectors
WithMetaInjectors adds Metadata injectors to the publisher. Injectors are run in the order they are given.
publisher.go:39
↓ 3 callersFunctionWithMetadataFilter
(filter *MetadataFilter)
consumer.go:252
↓ 3 callersFunctionisErrorCode
var ( _ pgError = (*pgconn.PgError)(nil) _ pgError = (*pq.Error)(nil) _ legacyPGError = (*pq.Error)(nil) )
errors.go:33
↓ 2 callersMethodGet
(k byte)
errors.go:24
↓ 2 callersMethodHasParam
HasParam checks whether the Builder has a parameter of the given name.
internal/query/query_builder.go:33
↓ 2 callersMethodbuildArgs
(ctx context.Context, msgs []*MessageOutgoing)
publisher.go:141
↓ 2 callersFunctionbuildInsertQuery
(queue string, msgCount int)
publisher.go:117
↓ 2 callersMethoddiscardMessage
(exec execer, msgID pgtype.UUID)
consumer.go:759
↓ 2 callersFunctiongenerateCreateTableQuery
(queueName string)
validator_test.go:205
↓ 2 callersFunctiongenerateInvalidQueueQuery
(queueName string)
validator_test.go:192
↓ 2 callersMethodgenerateQuery
()
consumer.go:375
↓ 2 callersFunctiongetParams
getParams extracts tags formatted as :tagName from the provided line, ignoring type casting patterns like ::interval.
internal/query/query_builder.go:61
↓ 2 callersFunctionisJSONObject
(b json.RawMessage)
consumer.go:702
↓ 2 callersFunctionprepareCtxTimeout
()
consumer.go:477
↓ 1 callersMethodBuild
Build returns the query string with the parameters replaced by the values from the provided map.
internal/query/query_builder.go:47
↓ 1 callersMethodHandleMessage
(context.Context, *MessageIncoming)
consumer.go:66
↓ 1 callersMethodLastAttempt
LastAttempt returns true if the message is consumed for the last time according to maxConsumedCount settings. If the Consumer is not configured to lim
message.go:105
↓ 1 callersFunctionNewBuilder
()
internal/query/query_builder.go:19
↓ 1 callersFunctionNewConsumerExt
NewConsumer creates Consumer with proper settings, using sqlx.DB (until refactored to use pgx directly)
consumer.go:269
↓ 1 callersFunctionNewPublisherExt
NewPublisher initializes the publisher with given *sqlx.DB client
publisher.go:59
↓ 1 callersMethodSQLState
()
errors.go:15
↓ 1 callersFunctionWithAckTimeout
WithAckTimeout sets the timeout for updating the message status when message is processed.
consumer.go:167
↓ 1 callersFunctionWithMaxConsumeCount
WithMaxConsumeCount sets the maximal number of times a message can be consumed before it is ignored. Unhandled SIGKILL or uncaught panic, OOM error et
consumer.go:216
↓ 1 callersFunctionWithMessageProcessingReserveDuration
WithMessageProcessingReserveDuration sets the duration for which the message is reserved for handling result state.
consumer.go:189
↓ 1 callersMethodack
ack positively acknowledges the message, and the message is marked as processed.
message.go:149
↓ 1 callersMethodackMessage
(exec execer, msgID pgtype.UUID)
consumer.go:715
↓ 1 callersFunctionbuildColumnListFromTags
buildColumnListFromTags dynamically constructs a list of column names based on the `db` struct tags of any given struct. It returns a slice of strings
message.go:183
↓ 1 callersFunctioncheckIndexData
(ctx context.Context, db *sqlx.DB, queueName string)
validator.go:125
↓ 1 callersMethodconsumeMessages
(ctx context.Context, query *query.Builder)
consumer.go:487
↓ 1 callersMethoddiscard
discard removes the message from the queue completely. It's like ack, but it also records the reason why the message was discarded.
message.go:171
↓ 1 callersMethoddiscardInvalidMsg
(ctx context.Context, id pgtype.UUID, err error)
consumer.go:639
↓ 1 callersFunctionensureUUIDExtension
(t *testing.T, db *sqlx.DB)
validator_test.go:170
↓ 1 callersFunctionensureUUIDExtension
(t *testing.T, db *sql.DB)
integtest/consumer_test.go:339
↓ 1 callersMethodfinishParsing
(pgMsg pgMessage)
consumer.go:654
↓ 1 callersFunctiongenerateCreateTablePartitionedQuery
(queueName string)
validator_test.go:243
↓ 1 callersFunctiongenerateCreateTableQueryCompositeIndex
(queueName string)
validator_test.go:224
↓ 1 callersFunctiongetColumnData
(ctx context.Context, db *sqlx.DB, queueName string)
validator.go:104
↓ 1 callersMethodhandleMessage
(ctx context.Context, msg *MessageIncoming)
consumer.go:417
↓ 1 callersMethodlogFields
(msg pgMessage, err error)
consumer.go:625
↓ 1 callersMethodnack
nack does not the negative acknowledge of the message. The message is returned to the queue after nack and may be processed again.
message.go:160
↓ 1 callersMethodnackMessage
(exec execer, msgID pgtype.UUID)
consumer.go:737
↓ 1 callersFunctionparseMetadata
(pgMsg pgMessage)
consumer.go:688
↓ 1 callersFunctionparsePayload
(pgMsg pgMessage)
consumer.go:678
↓ 1 callersMethodparseRow
(ctx context.Context, rows *sqlx.Rows)
consumer.go:597
↓ 1 callersFunctionprepareProcessMetric
(queueName string, meter metric.Meter)
consumer.go:294
↓ 1 callersMethodtryConsumeMessages
(ctx context.Context, query *query.Builder, limit int64)
consumer.go:523
↓ 1 callersMethodupdateLockedUntil
(db *sqlx.DB, id pgtype.UUID)
consumer.go:781
↓ 1 callersMethodverifyTable
(ctx context.Context)
consumer.go:357
MethodError
()
consumer.go:37
MethodError
()
errors_test.go:15
MethodError
()
errors_test.go:31
FunctionExampleConsumer
()
example_consumer_test.go:50
FunctionExampleNewConsumer
()
examples_test.go:16
FunctionExampleNewPublisher
()
examples_test.go:39
FunctionExamplePublisher
()
example_publisher_test.go:17
MethodFatal
()
errors_test.go:19
MethodGet
(k byte)
errors_test.go:23
MethodHandleMessage
HandleMessage calls self. It also implements MessageHandler interface.
consumer.go:73
MethodHandleMessage
(ctx context.Context, msg *pgq.MessageIncoming)
example_consumer_test.go:18
MethodHandleMessage
(ctx context.Context, _ *MessageIncoming)
integtest/consumer_test.go:357
MethodHandleMessage
(ctx context.Context, _ *MessageIncoming)
integtest/consumer_test.go:362
FunctionNewMessage
NewMessage creates new message that satisfies Message interface.
message.go:55
MethodPublish
Publish publishes the message.
publisher.go:68
MethodSQLState
()
errors_test.go:35
MethodSetTimeout
SetTimeout sets the message timeout. If the timeout is after the message deadline, it returns ErrInvalidDeadline. The timeout also affects the queue l
message.go:117
FunctionTestAcquireMaxFromSemaphore
(t *testing.T)
consumer_test.go:93
FunctionTestClient_buildArgs
(t *testing.T)
publisher_test.go:48
next →1–100 of 122, ranked by callers