Home
last modified time | relevance | path

Searched refs:feed (Results 1 – 25 of 68) sorted by relevance

123

/trunk/goproj/src/github.com/couchbase/indexing/secondary/projector/
Dfeed.go321 feed.shutdown(feed.opaque)
493 feed.shutdown(feed.opaque)
495 feed.projector.DelFeed(feed.topic)
763 feed.shutdown(feed.opaque)
788 feed.shutdown(feed.opaque)
1590 feed.projector.UpdateStats(feed.topic, feed)
1901 feed.kvaddr, feed.config, feed.async, opaque2, feed.collectionsAware)
2232 feed.logPrefix, feed.opaque, getTssAsStr(feed.actTss))
2244 feed.logPrefix, feed.opaque, getTssAsStr(feed.actTss), getTssAsStr(feed.rollTss))
2253 feed.logPrefix, feed.opaque, keyspaceId, getTssAsStr(feed.actTss), getTssAsStr(feed.reqTss))
[all …]
Dprojector.go432 return feed, ok
439 return feed, nil
465 p.topics[topic] = feed
466 opaque := feed.GetOpaque()
482 opaque := feed.GetOpaque()
647 if feed == nil {
660 p.UpdateStats(topic, feed)
661 p.AddFeed(topic, feed)
724 p.UpdateStats(topic, feed)
785 p.UpdateStats(topic, feed)
[all …]
Dkvdata.go46 feed *Feed member
214 feed *Feed,
229 feed: feed,
231 topic: feed.topic,
274 kvdata.logPrefix = fmt.Sprintf(fmsg, keyspaceId, feed.cluster, feed.topic)
476 kvdata.feed.PostFinKVdata(kvdata.keyspaceId, kvdata.uuid)
665 kvdata.feed.PostStreamRequest(kvdata.keyspaceId, m, kvdata.uuid)
682 kvdata.feed.PostStreamEnd(kvdata.keyspaceId, m, kvdata.uuid)
758 feed *Feed, bucket, keyspaceId string, config c.Config,
764 workers[i] = NewVbucketWorker(i, feed, bucket,
[all …]
Ddcp_seqno_local.go348 for _, feed := range kvfeeds {
350 go func(index int, feed *kvConn) {
352 kv_seqnos_node[index] = feed.seqsbuf
354 errors[index] = couchbase.GetSeqs(feed.mc, kv_seqnos_node[index], feed.tmpbuf)
357 err := tryEnableCollection(feed.mc)
362 errors[index] = couchbase.GetCollectionSeqs(feed.mc, kv_seqnos_node[index], feed.tmpbuf, cid)
365 }(i, feed)
/trunk/goproj/src/github.com/couchbase/go-couchbase/
Dupr.go115 feed := &UprFeed{
132 go feed.run()
134 feed.C = feed.output
135 return feed, nil
290 …singleFeed, err = serverConn.StartUprFeed(name, feed.sequence, feed.dcp_buffer_size, feed.data_cha…
305 feed.wg.Add(1)
320 feed.wg.Done()
323 if feed.outputClosed != true && feed.closing != true {
383 case <-feed.quit:
390 close(feed.quit)
[all …]
Dtap.go36 feed := &TapFeed{
43 go feed.run()
45 feed.C = feed.output
46 return feed, nil
77 case <-feed.quit:
96 feed.nodeFeeds = append(feed.nodeFeeds, singleFeed)
98 feed.wg.Add(1)
117 case <-feed.quit:
133 case <-feed.quit:
139 close(feed.quit)
[all …]
/trunk/goproj/src/github.com/couchbase/gomemcached/client/
Dupr_feed.go351 feed.Close()
395 feed := &UprFeed{
410 return feed, nil
485 if !feed.isOpen() {
510 mc := feed.conn
643 feed.setOpen()
755 return feed.C
759 return feed.Error
796 feed.C = ch
1062 feed.Close()
[all …]
Dtap_feed.go247 feed := &TapFeed{
251 go mc.runFeed(ch, feed)
252 return feed, nil
260 func (mc *Client) runFeed(ch chan TapEvent, feed *TapFeed) {
277 feed.Error = err
287 feed.Error = fmt.Errorf("tap connection failed: %s", pkt.Body)
299 case <-feed.closer:
308 feed.Error = err
331 func (feed *TapFeed) Close() {
332 close(feed.closer)
/trunk/goproj/src/github.com/couchbase/indexing/secondary/dcp/transport/client/
Ddcp_feed.go127 feed.conn = mc
136 feed.isIncrBuild = feed.collectionsAware
143 go feed.genServer(opaque, feed.reqch, feed.finch, rcvch, config)
144 go feed.doReceive(rcvch, feed.finch, mc)
252 feed.sendStreamEnd(feed.outch)
336 feed.sendStreamEnd(feed.outch)
344 feed.sendStreamEnd(feed.outch)
350 feed.sendStreamEnd(feed.outch)
981 feed.isIncrBuild = feed.isIncrBuild && false
987 feed.isIncrBuild = feed.isIncrBuild && true
[all …]
Dtap_feed.go247 feed := &TapFeed{
251 go mc.runFeed(ch, feed)
252 return feed, nil
260 func (mc *Client) runFeed(ch chan TapEvent, feed *TapFeed) {
277 feed.Error = err
287 feed.Error = fmt.Errorf("tap connection failed: %s", pkt.Body)
299 case <-feed.closer:
308 feed.Error = err
331 func (feed *TapFeed) Close() {
332 close(feed.closer)
/trunk/goproj/src/github.com/couchbase/indexing/secondary/dcp/
Dupr.go178 feed.C = feed.output
184 go feed.genServer(feed.reqch, opaque)
185 return feed, nil
216 resp, err := failsafeOp(feed.reqch, respch, cmd, feed.finch)
226 resp, err := failsafeOp(feed.reqch, respch, cmd, feed.finch)
240 resp, err := failsafeOp(feed.reqch, respch, cmd, feed.finch)
251 resp, err := failsafeOp(feed.reqch, respch, cmd, feed.finch)
258 resp, err := failsafeOp(feed.reqch, respch, cmd, feed.finch)
384 feedname, feed.sequence, flags, feed.output, opaque, feed.reqch, config)
428 feedname, feed.sequence, flags, feed.output, opaque, feed.reqch, config)
[all …]
Dtap.go35 feed := &TapFeed{
42 go feed.run()
44 feed.C = feed.output
45 return feed, nil
60 case <-feed.quit:
76 case <-feed.quit:
95 feed.nodeFeeds = append(feed.nodeFeeds, singleFeed)
114 case <-feed.quit:
130 case <-feed.quit:
136 close(feed.quit)
[all …]
/trunk/cbft/
H A Dpindex_bleve_rollback_test.go28 feed.OpaqueSet(partition, buf)
93 var feed *cbgt.PrimaryFeed
155 feed, ok = f.(*cbgt.PrimaryFeed)
170 err = feed.SnapshotStart(partition, snapStart, snapEnd)
174 setVBucketFailoverLog(feed, partition)
190 err = feed.DataUpdate(partition, key, seq, val,
211 err = feed.DataUpdate(partition, key, seq, val,
231 err = feed.SnapshotStart(partition, snapStart, snapEnd)
251 err = feed.DataUpdate(partition, key, seq, val,
293 err = feed.DataUpdate(partition, key, seq, val,
[all …]
H A Drest_test.go839 Desc: "re-create a deleted index with nil feed",
974 t.Errorf("expected the 1 feed to be a PrimaryFeed")
1616 Desc: "delete idx0 with dest feed",
2024 t.Errorf("expected 1 feed, got feeds: %+v", feeds)
2043 t.Errorf("expected 1 feed, got feeds: %+v", feeds)
2308 t.Errorf("expected to be %d feed, got feeds: %+v",
2865 Desc: "create an index with nil feed",
3099 Desc: "create another index with primary feed",
3114 Desc: "create another index with primary feed",
3129 Desc: "feed stats on a multi-index manager",
[all …]
/trunk/goproj/src/github.com/couchbase/go-couchbase/examples/upr_restart/
Drestart.go26 feed, err := bucket.StartUprFeed("index" /*name*/, 0)
32 err := feed.UprRequestStream(
40 vbseqNo := receiveMutations(feed, 20000)
63 feed, err = bucket.StartUprFeed("index" /*name*/, 0)
71 err := feed.UprRequestStream(
86 case f = <-feed.C:
106 feed.Close()
119 func receiveMutations(feed *couchbase.UprFeed, breakAfter int) [][2]uint64 {
127 case e = <-feed.C:
/trunk/goproj/src/github.com/couchbase/indexing/secondary/dcp/examples/upr_restart/
Drestart.go67 feed, err := bucket.StartDcpFeed(name, 0, 0xABCD, dcpConfig)
73 err := feed.DcpRequestStream(
81 vbseqNo := receiveMutations(feed, 20000)
104 feed, err = bucket.StartDcpFeed(name, 0, 0xABCD, dcpConfig)
112 err := feed.DcpRequestStream(
127 case f := <-feed.C:
145 feed.Close()
158 func receiveMutations(feed *couchbase.DcpFeed, breakAfter int) [][2]uint64 {
165 case e := <-feed.C:
/trunk/cbgt/
H A Dfeed_tap.go47 feed, err := NewTAPFeed(feedName, indexName, server, poolName,
55 err = feed.Start()
60 err = mgr.registerFeed(feed)
62 feed.Close()
156 progress, err := t.feed()
168 func (t *TAPFeed) feed() (int, error) { func
201 feed, err := bucket.StartTapFeed(&args)
205 defer feed.Close()
225 case req, alive := <-feed.C:
H A Dfeed_dcp_gocbcore.go413 err = mgr.registerFeed(feed)
417 return feed.onError(false, err)
420 err = feed.Start()
422 return feed.onError(true,
565 feed := &GocbcoreDCPFeed{
595 feed.scope = "_default"
598 feed.scope = params.Scope
617 feed.agent, err = FetchDCPAgent(feed.bucketName, feed.bucketUUID,
620 return nil, feed.onSetupError(
626 name, indexName, feed.servers, feed.bucketName, feed.bucketUUID)
[all …]
H A Dfeed_dcp_gocouchbase.go102 feed, err := NewDCPFeed(feedName, indexName, server, poolName,
111 err = feed.Start()
120 feed.name, cbdatasource.UpdateSecuritySettings, status)
121 return feed.bds.Kick(cbdatasource.UpdateSecuritySettings)
124 RegisterConfigRefreshCallback("DCPFeed_"+feed.name,
127 err = mgr.registerFeed(feed)
129 feed.Close()
259 feed := &DCPFeed{
277 feed.bds, err = cbdatasource.NewBucketDataSource(
279 vbucketIds, auth, feed, options)
[all …]
/trunk/goproj/src/github.com/couchbase/indexing/secondary/dcp/examples/upr_feed/
Dfeed.go110 feed, err := bucket.StartDcpFeed(name, 0, 0xABCD, dcpConfig)
119 err := feed.DcpRequestStream(
128 events(feed, 2000)
138 if err := feed.DcpCloseStream(uint16(vbno), opaque); err != nil {
145 events(feed, 100000)
147 feed.Close()
150 func events(feed *couchbase.DcpFeed, timeoutMs int) {
163 case e := <-feed.C:
/trunk/libcouchbase/tests/rdb/
H A Dt_refs.cc11 ior->feed("12345678"); in TEST_F()
44 ior->feed("1234567890A"); in TEST_F()
80 ior->feed("123456789"); in TEST_F()
95 ior->feed("123456789"); in TEST_F()
104 ior->feed("123456789"); in TEST_F()
/trunk/goproj/src/github.com/couchbase/go-couchbase/examples/upr_feed/
Dfeed.go91 feed, err := bucket.StartUprFeed(name, 0)
99 err := feed.UprRequestStream(
114 case e, ok := <-feed.C:
141 err := feed.UprCloseStream(
148 feed.Close()
/trunk/goproj/src/github.com/couchbase/indexing/secondary/common/
Ddcp_seqno.go955 for kvaddr, feed := range kvfeeds {
957 go func(kvaddress string, feed *kvConn) {
961 err := couchbase.GetSeqs(feed.mc, feed.seqsbuf, feed.tmpbuf)
970 err := tryEnableCollection(feed.mc)
977 err = couchbase.GetCollectionSeqs(feed.mc, feed.seqsbuf, feed.tmpbuf, cid)
990 }(kvaddr, feed)
1350 for kvaddr, feed := range kvfeeds {
1357 feed.seqsbuf, feed.tmpbuf)
1366 err := tryEnableCollection(feed.mc)
1374 feed.seqsbuf, feed.tmpbuf, cid)
[all …]
/trunk/goproj/src/github.com/couchbase/indexing/secondary/dcp/transport/client/example/
Dtap_example.go45 feed, err := client.StartTapFeed(args)
49 for op := range feed.C {
59 log.Printf("Tap feed closed; err = %v.", feed.Error)
/trunk/goproj/src/github.com/couchbase/gomemcached/client/example/
Dtap_example.go45 feed, err := client.StartTapFeed(args)
49 for op := range feed.C {
59 log.Printf("Tap feed closed; err = %v.", feed.Error)

123