MCPcopy Create free account

hub / github.com/apache/pulsar-client-go / functions

Functions5,094 in github.com/apache/pulsar-client-go

↓ 3 callersFunctionsortMessageIDs
(msgIDs []messageID)
pulsar/negative_acks_tracker_test.go:72
↓ 3 callersFunctionsubscriber
(c *client, topics []string, opts ConsumerOptions, ch chan ConsumerMessage, dlq *dlqRouter, rlq *retryRouter)
pulsar/consumer_regex.go:478
↓ 3 callersMethodsuccess
()
pulsar/consumer_test.go:6249
↓ 3 callersFunctionvalidateTopicNames
(topics ...string)
pulsar/helper.go:56
↓ 3 callersFunctionverifyLogOutput
(t *testing.T, logOutput, expectedLevel, expectedMessage string, expectedFields ...Fields)
pulsar/log/wrapper_slog_test.go:169
↓ 3 callersMethodverifyOpen
()
pulsar/transaction_impl.go:223
↓ 3 callersMethodwaitWithContext
waitWithContext Same as wait() call, but the end condition can also be controlled through the context. It blocks until either a broadcast occurs or th
pulsar/internal/channel_cond.go:51
↓ 3 callersMethodwriteData
(buffer internal.Buffer, sequenceID uint64, callbacks []interface{})
pulsar/producer_partition.go:904
↓ 2 callersMethodActiveConsumerChanged
(isActive bool)
pulsar/internal/connection.go:97
↓ 2 callersMethodBeforeConsume
BeforeConsume This is called just before the message is send to Consumer's ConsumerMessage channel.
pulsar/consumer_interceptor.go:22
↓ 2 callersMethodBeforeSend
BeforeSend This is called before send the message to the brokers. This method is allowed to modify the message.
pulsar/producer_interceptor.go:23
↓ 2 callersFunctionCheckName
(name string)
pulsaradmin/pkg/utils/namespace_name.go:86
↓ 2 callersMethodClose
()
pulsar/consumer_partition.go:1030
↓ 2 callersMethodClose
()
pulsar/internal/client_handlers.go:52
↓ 2 callersFunctionConvertGetAllSchemasResponseToSchemaInfosWithVersion
( tn *TopicName, response GetAllSchemasResponse, )
pulsaradmin/pkg/utils/schema_util.go:98
↓ 2 callersFunctionCrc32cCheckSum
Crc32cCheckSum handles computing the checksum.
pulsar/internal/checksum.go:34
↓ 2 callersMethodCreateNamespace
CreateNamespace creates a new empty namespace with no policies attached
pulsaradmin/pkg/admin/namespace.go:50
↓ 2 callersMethodCreateNsWithBundlesDataWithContext
( ctx context.Context, namespace string, bundleData *utils.BundlesData, )
pulsaradmin/pkg/admin/namespace.go:846
↓ 2 callersMethodCreateSchemaByPayloadWithContext
( ctx context.Context, topic string, schemaPayload utils.PostSchemaPayload, )
pulsaradmin/pkg/admin/schema.go:213
↓ 2 callersMethodCreateWithPropertiesWithContext
( ctx context.Context, topic utils.TopicName, partitions int, meta map[string]string, )
pulsaradmin/pkg/admin/topic.go:1150
↓ 2 callersFunctionDLQWithProducerOptions
(t *testing.T, prodOpt *ProducerOptions)
pulsar/consumer_test.go:1941
↓ 2 callersMethodDeleteWithContext
DeleteWithContext deletes an existing cluster
pulsaradmin/pkg/admin/cluster.go:50
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:80
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:142
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:271
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:330
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:448
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:504
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:560
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:673
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:735
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:791
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:847
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:906
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:962
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:1030
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:1086
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:1145
↓ 2 callersMethodDescriptor
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:1370
↓ 2 callersMethodEncrypt
([]byte, *pb.MessageMetadata)
pulsar/internal/crypto/encryptor.go:26
↓ 2 callersFunctionExtractSpanContextFromProducerMessage
(message *pulsar.ProducerMessage)
pulsar/internal/pulsartracing/message_carrier_util.go:40
↓ 2 callersMethodFloat64
(buf []byte)
pulsar/primitiveSerDe.go:101
↓ 2 callersMethodFlushBatches
FlushBatches all the messages buffered in multiple batches and wait until all messages have been successfully persisted.
pulsar/internal/batch_builder.go:60
↓ 2 callersMethodGetAckSet
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:1549
↓ 2 callersMethodGetAssignedBrokerServiceUrl
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:4604
↓ 2 callersMethodGetAssignedBrokerServiceUrlTls
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:4611
↓ 2 callersFunctionGetAuthProvider
(config *config.Config)
pulsaradmin/pkg/admin/auth/provider.go:57
↓ 2 callersMethodGetBrokerAddress
(brokerServiceURL string, proxyThroughServiceURL bool)
pulsar/internal/lookup_service.go:76
↓ 2 callersMethodGetBrokerAddress
(brokerServiceURL string, proxyThroughServiceURL bool)
pulsar/internal/lookup_service.go:137
↓ 2 callersMethodGetBundleRange
GetBundleRange returns a bundle range of a topic
pulsaradmin/pkg/admin/topic.go:225
↓ 2 callersFunctionGetCompressionProvider
( compressionType pb.CompressionType, level compression.Level, )
pulsar/internal/batch_builder.go:308
↓ 2 callersMethodGetConnections
GetConnections get all connections in the pool.
pulsar/internal/connection_pool.go:38
↓ 2 callersFunctionGetConnectionsCount
(p *ConnectionPool)
pulsar/internal/helper.go:28
↓ 2 callersMethodGetEncryptionContext
GetEncryptionContext returns the ecryption context of the message. It will be used by the application to parse the undecrypted message.
pulsar/message.go:140
↓ 2 callersMethodGetLastMessageId
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:5384
↓ 2 callersMethodGetMaxTopicsPerNamespace
GetMaxTopicsPerNamespace returns the maxTopicsPerNamespace for a namespace.
pulsaradmin/pkg/admin/namespace.go:364
↓ 2 callersMethodGetMessageID
GetMessageID returns the message Id by timestamp(ms) of a topic @param topic topicName struct @param timestamp absolute timestamp (in ms)
pulsaradmin/pkg/admin/topic.go:242
↓ 2 callersMethodGetMessagesByIDWithContext
( ctx context.Context, topic utils.TopicName, ledgerID, entryID int64, )
pulsaradmin/pkg/admin/subscription.go:347
↓ 2 callersMethodGetNullValue
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:2053
↓ 2 callersMethodGetOrderingKey
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:2011
↓ 2 callersFunctionGetPackageNameWithComponents
(packageType PackageType, tenant, namespace, name, version string)
pulsaradmin/pkg/utils/package_name.go:42
↓ 2 callersMethodGetPayloadSize
()
pulsaradmin/pkg/utils/message.go:73
↓ 2 callersMethodGetProducerReady
()
pulsar/internal/pulsar_proto/PulsarApi.pb.go:4901
↓ 2 callersMethodGetSchemaCompatibilityStrategyAppliedWithContext
( ctx context.Context, topic utils.TopicName, applied bool, )
pulsaradmin/pkg/admin/topic.go:2361
↓ 2 callersMethodGetState
GetState Get the state of the transaction.
pulsar/transaction.go:66
↓ 2 callersMethodGetStatsWithOptionWithContext
GetStatsWithOptionWithContext returns the stats for the topic All the rates are computed over a 1-minute window and are relative the last completed 1
pulsaradmin/pkg/admin/topic.go:284
↓ 2 callersMethodGetSubPermissions
GetSubPermissions returns subscription permissions on a namespace
pulsaradmin/pkg/admin/namespace.go:532
↓ 2 callersMethodGetSubscriptionExpirationTime
GetSubscriptionExpirationTime gets the subscription expiration time on a namespace. Returns -1 if not set
pulsaradmin/pkg/admin/namespace.go:721
↓ 2 callersMethodGetTLSCertificate
()
pulsar/auth/tls.go:64
↓ 2 callersMethodGetTokens
(identifier string)
oauth2/config_tokenprovider.go:23
↓ 2 callersMethodGetTopicAutoCreation
GetTopicAutoCreation returns the topic auto-creation config for a namespace. Returns nil if the topic auto-creation config is not configured at the na
pulsaradmin/pkg/admin/namespace.go:151
↓ 2 callersMethodGetVersionByPayloadWithContext
( ctx context.Context, topic string, schemaPayload utils.PostSchemaPayload, )
pulsaradmin/pkg/admin/schema.go:259
↓ 2 callersMethodHasURL
()
pulsar/producer_partition.go:416
↓ 2 callersMethodHealthCheckWithTopicVersion
HealthCheckWithTopicVersion runs a health check on the broker
pulsaradmin/pkg/admin/brokers.go:100
↓ 2 callersMethodHealthCheckWithTopicVersionWithContext
(ctx context.Context, topicVersion utils.TopicVersion)
pulsaradmin/pkg/admin/brokers.go:260
↓ 2 callersMethodIndex
Index returns index from broker entry metadata, or empty if the feature is not enabled in the broker.
pulsar/message.go:144
↓ 2 callersFunctionInjectConsumerMessageSpanContext
(ctx context.Context, message pulsar.ConsumerMessage)
pulsar/internal/pulsartracing/message_carrier_util.go:64
↓ 2 callersFunctionInjectProducerMessageSpanContext
(ctx context.Context, message *pulsar.ProducerMessage)
pulsar/internal/pulsartracing/message_carrier_util.go:28
↓ 2 callersMethodInvalidateToken
InvalidateToken is called when the token is rejected by the resource server.
oauth2/cache/cache.go:35
↓ 2 callersMethodIsHTTP
()
pulsar/internal/service_uri.go:85
↓ 2 callersMethodIsMaxBackoffReached
IsMaxBackoffReached evaluates if the max number of retries is reached
pulsar/backoff/backoff.go:32
↓ 2 callersMethodIsNullValue
IsNullValue reports whether the message was published as a null-value (tombstone) message, i.e. with MessageMetadata.null_value set. For such messages
pulsar/message.go:97
↓ 2 callersMethodIsProxied
()
pulsar/internal/connection.go:91
↓ 2 callersFunctionMakeHTTPPath
(apiVersion string, componentPath string)
pulsaradmin/pkg/utils/utils.go:26
↓ 2 callersMethodMakeRequestWithURLWithContext
( ctx context.Context, method string, urlOpt *url.URL, )
pulsaradmin/pkg/rest/client.go:117
↓ 2 callersFunctionNewAuthDisabled
NewAuthDisabled return a interface of Provider
pulsar/auth/disabled.go:28
↓ 2 callersFunctionNewAuthenticationAthenzWithParams
(params map[string]string)
pulsar/auth/athenz.go:70
↓ 2 callersFunctionNewAuthenticationFromTLSCertSupplier
NewAuthenticationFromTLSCertSupplier Create new Authentication provider with specified TLS certificate supplier
pulsar/client.go:66
↓ 2 callersFunctionNewAuthenticationOAuth2WithFlow
( issuer oauth2.Issuer, flowOptions oauth2.ClientCredentialsFlowOptions)
pulsaradmin/pkg/admin/auth/oauth2.go:60
↓ 2 callersFunctionNewAuthenticationTLS
NewAuthenticationTLS initialize the authentication provider
pulsar/auth/tls.go:41
↓ 2 callersFunctionNewAuthenticationTLS
NewAuthenticationTLS initialize the authentication provider
pulsaradmin/pkg/admin/auth/tls.go:44
↓ 2 callersFunctionNewAuthenticationToken
NewAuthenticationToken returns a token auth provider that will use the specified token to talk with Pulsar brokers
pulsar/auth/token.go:47
↓ 2 callersFunctionNewAuthenticationTokenFromFile
NewAuthenticationTokenFromFile return a interface of a Provider with a string token file path.
pulsar/auth/token.go:68
↓ 2 callersFunctionNewAuthenticationTokenFromSupplier
NewAuthenticationTokenFromSupplier returns a token auth provider that gets the token data from a user supplied function. The function is invoked each
pulsar/client.go:51
↓ 2 callersFunctionNewBufferPool
()
pulsar/internal/buffer.go:97
↓ 2 callersFunctionNewBundlesDataWithNumBundles
(numBundles int)
pulsaradmin/pkg/utils/bundles_data.go:32
↓ 2 callersFunctionNewClient
()
perf/pulsar-perf-go.go:53
↓ 2 callersFunctionNewDefaultBackoffWithInitialBackOff
(backoff time.Duration)
pulsar/backoff/backoff.go:50
↓ 2 callersFunctionNewDefaultTokenCache
(audience string, flow *oauth2.ClientCredentialsFlow)
oauth2/cache/cache.go:53
↓ 2 callersFunctionNewMessageReader
(headersAndPayload Buffer)
pulsar/internal/commands.go:54
← previousnext →601–700 of 5,094, ranked by callers