Search Options

Display Count
Sort
Preferred Language
Advanced Search

Results 11 - 19 of 19 for queue_dir (0.06 seconds)

The search processing time has exceeded the limit. The displayed results may be partial.

  1. internal/event/target/nats.go

    func NewNATSTarget(id string, args NATSArgs, loggerOnce logger.LogOnce) (*NATSTarget, error) {
    	var queueStore store.Store[event.Event]
    	if args.QueueDir != "" {
    		queueDir := filepath.Join(args.QueueDir, storePrefix+"-nats-"+id)
    		queueStore = store.NewQueueStore[event.Event](queueDir, args.QueueLimit, event.StoreExtension)
    		if err := queueStore.Open(); err != nil {
    Created: Sun Apr 05 19:28:12 GMT 2026
    - Last Modified: Sun Apr 27 04:30:57 GMT 2025
    - 13.5K bytes
    - Click Count (0)
  2. internal/event/target/elasticsearch.go

    	var queueStore store.Store[event.Event]
    	if args.QueueDir != "" {
    		queueDir := filepath.Join(args.QueueDir, storePrefix+"-elasticsearch-"+id)
    		queueStore = store.NewQueueStore[event.Event](queueDir, args.QueueLimit, event.StoreExtension)
    		if err := queueStore.Open(); err != nil {
    Created: Sun Apr 05 19:28:12 GMT 2026
    - Last Modified: Sun Sep 28 20:59:21 GMT 2025
    - 15K bytes
    - Click Count (0)
  3. internal/config/notify/legacy.go

    				}
    				return strings.Join(brokers, config.ValueSeparator)
    			}(),
    		},
    		config.KV{
    			Key:   target.KafkaTopic,
    			Value: cfg.Topic,
    		},
    		config.KV{
    			Key:   target.KafkaQueueDir,
    			Value: cfg.QueueDir,
    		},
    		config.KV{
    			Key:   target.KafkaClientTLSCert,
    			Value: cfg.TLS.ClientTLSCert,
    		},
    		config.KV{
    			Key:   target.KafkaClientTLSKey,
    			Value: cfg.TLS.ClientTLSKey,
    		},
    		config.KV{
    Created: Sun Apr 05 19:28:12 GMT 2026
    - Last Modified: Sun Apr 27 04:30:57 GMT 2025
    - 13.3K bytes
    - Click Count (0)
  4. internal/logger/target/kafka/kafka.go

    		User      string `json:"username"`
    		Password  string `json:"password"`
    		Mechanism string `json:"mechanism"`
    	} `json:"sasl"`
    	// Queue store
    	QueueSize int    `json:"queueSize"`
    	QueueDir  string `json:"queueDir"`
    
    	// Custom logger
    	LogOnce func(ctx context.Context, err error, id string, errKind ...any) `json:"-"`
    }
    
    // Target - Kafka target.
    type Target struct {
    	status int32
    
    	totalMessages  int64
    Created: Sun Apr 05 19:28:12 GMT 2026
    - Last Modified: Sun Sep 28 20:59:21 GMT 2025
    - 10.2K bytes
    - Click Count (0)
  5. internal/logger/help.go

    		},
    		config.HelpKV{
    			Key:         QueueSize,
    			Description: "configure channel queue size for webhook targets",
    			Optional:    true,
    			Type:        "number",
    		},
    		config.HelpKV{
    			Key:         QueueDir,
    			Description: `staging dir for undelivered logger messages e.g. '/home/logger-events'`,
    			Optional:    true,
    			Type:        "string",
    		},
    		config.HelpKV{
    			Key:         Proxy,
    Created: Sun Apr 05 19:28:12 GMT 2026
    - Last Modified: Wed Sep 11 22:20:42 GMT 2024
    - 7.4K bytes
    - Click Count (0)
  6. internal/logger/target/http/http.go

    	ClientCert  string            `json:"clientCert"`
    	ClientKey   string            `json:"clientKey"`
    	BatchSize   int               `json:"batchSize"`
    	QueueSize   int               `json:"queueSize"`
    	QueueDir    string            `json:"queueDir"`
    	MaxRetry    int               `json:"maxRetry"`
    	RetryIntvl  time.Duration     `json:"retryInterval"`
    	Proxy       string            `json:"string"`
    	Transport   http.RoundTripper `json:"-"`
    Created: Sun Apr 05 19:28:12 GMT 2026
    - Last Modified: Fri Aug 29 02:39:48 GMT 2025
    - 15.6K bytes
    - Click Count (0)
  7. internal/store/batch_test.go

    	"time"
    )
    
    func TestBatchCommit(t *testing.T) {
    	defer func() {
    		if err := tearDownQueueStore(); err != nil {
    			t.Fatalf("Failed to tear down store; %v", err)
    		}
    	}()
    	store, err := setUpQueueStore(queueDir, 100)
    	if err != nil {
    		t.Fatalf("Failed to create a queue store; %v", err)
    	}
    
    	var limit uint32 = 100
    
    	batch := NewBatch[TestItem](BatchConfig[TestItem]{
    		Limit:         limit,
    		Store:         store,
    Created: Sun Apr 05 19:28:12 GMT 2026
    - Last Modified: Fri Aug 29 02:39:48 GMT 2025
    - 5.6K bytes
    - Click Count (0)
  8. internal/config/notify/parse.go

    		if err != nil {
    			return nil, err
    		}
    		kafkaArgs := target.KafkaArgs{
    			Enable:             enabled,
    			Brokers:            brokers,
    			Topic:              env.Get(topicEnv, kv.Get(target.KafkaTopic)),
    			QueueDir:           env.Get(queueDirEnv, kv.Get(target.KafkaQueueDir)),
    			QueueLimit:         queueLimit,
    			Version:            env.Get(versionEnv, kv.Get(target.KafkaVersion)),
    			BatchSize:          uint32(batchSize),
    Created: Sun Apr 05 19:28:12 GMT 2026
    - Last Modified: Fri Aug 29 02:39:48 GMT 2025
    - 47.5K bytes
    - Click Count (0)
  9. cmd/config-versions.go

    // ConsoleLogger is introduced to workaround the dependency about logrus
    type ConsoleLogger struct {
    	Enable bool `json:"enable"`
    }
    
    // serverConfigV33 is just like version '32', removes clientID from NATS and MQTT, and adds queueDir, queueLimit in all notification targets.
    type serverConfigV33 struct {
    	quick.Config `json:"-"` // ignore interfaces
    
    	Version string `json:"version"`
    
    	// S3 API configuration.
    Created: Sun Apr 05 19:28:12 GMT 2026
    - Last Modified: Fri May 24 23:05:23 GMT 2024
    - 2.5K bytes
    - Click Count (0)
Back to Top