Compare commits

...

7 Commits

Author SHA1 Message Date
mr
d985d8339a Change of state Conn Management 2026-02-05 16:17:33 +01:00
mr
ea14ad3933 Closure On change of state 2026-02-05 16:17:14 +01:00
mr
2e31df89c2 oc-discovery + auto create peer 2026-02-05 15:47:29 +01:00
mr
425cbdfe7d stream address 2026-02-05 15:36:22 +01:00
mr
8ee5b84e21 publish-registry 2026-02-05 12:14:02 +01:00
mr
552bb17e2b Connectivity ok 2026-02-05 11:23:11 +01:00
mr
88e29073a2 dockerfile default 2026-02-05 09:31:51 +01:00
15 changed files with 108 additions and 81 deletions

View File

@@ -21,9 +21,6 @@ RUN go mod download
FROM golang:alpine AS builder FROM golang:alpine AS builder
ARG CONF_NUM ARG CONF_NUM
# Fail fast if CONF_NUM missing
RUN test -n "$CONF_NUM"
RUN apk add --no-cache git RUN apk add --no-cache git
WORKDIR /oc-discovery WORKDIR /oc-discovery
@@ -55,13 +52,13 @@ WORKDIR /app
RUN mkdir ./pem RUN mkdir ./pem
COPY --from=builder /app/extracted/pem/private${CONF_NUM}.pem ./pem/private.pem COPY --from=builder /app/extracted/pem/private${CONF_NUM:-1}.pem ./pem/private.pem
COPY --from=builder /app/extracted/psk ./psk COPY --from=builder /app/extracted/psk ./psk
COPY --from=builder /app/extracted/pem/public${CONF_NUM}.pem ./pem/public.pem COPY --from=builder /app/extracted/pem/public${CONF_NUM:-1}.pem ./pem/public.pem
COPY --from=builder /app/extracted/oc-discovery /usr/bin/oc-discovery COPY --from=builder /app/extracted/oc-discovery /usr/bin/oc-discovery
COPY --from=builder /app/extracted/docker_discovery${CONF_NUM}.json /etc/oc/discovery.json COPY --from=builder /app/extracted/docker_discovery${CONF_NUM:-1}.json /etc/oc/discovery.json
EXPOSE 400${CONF_NUM} EXPOSE 400${CONF_NUM:-1}
ENTRYPOINT ["oc-discovery"] ENTRYPOINT ["oc-discovery"]

View File

@@ -10,15 +10,17 @@ clean:
rm -rf oc-discovery rm -rf oc-discovery
docker: docker:
DOCKER_BUILDKIT=1 docker build -t oc/oc-discovery:0.0.1 -f Dockerfile . DOCKER_BUILDKIT=1 docker build -t oc-discovery -f Dockerfile .
docker tag oc/oc-discovery:0.0.1 oc/oc-discovery:latest docker tag oc-discovery opencloudregistry/oc-discovery:latest
publish-kind: publish-kind:
kind load docker-image oc/oc-discovery:0.0.1 --name opencloud kind load docker-image oc/oc-discovery:0.0.1 --name opencloud
publish-registry: publish-registry:
@echo "TODO" docker push opencloudregistry/oc-discovery:latest
all: docker publish-kind publish-registry all: docker publish-kind
ci: docker publish-registry
.PHONY: build run clean docker publish-kind publish-registry .PHONY: build run clean docker publish-kind publish-registry

View File

@@ -159,20 +159,16 @@ func (s *LongLivedPubSubService) SubscribeToSearch(ps *pubsub.PubSub, f *func(co
func SubscribeEvents[T interface{}](s *LongLivedPubSubService, func SubscribeEvents[T interface{}](s *LongLivedPubSubService,
ctx context.Context, proto string, timeout int, f func(context.Context, T, string), ctx context.Context, proto string, timeout int, f func(context.Context, T, string),
) error { ) error {
s.PubsubMu.Lock()
if s.LongLivedPubSubs[proto] == nil { if s.LongLivedPubSubs[proto] == nil {
s.PubsubMu.Unlock()
return errors.New("no protocol subscribed in pubsub") return errors.New("no protocol subscribed in pubsub")
} }
topic := s.LongLivedPubSubs[proto] topic := s.LongLivedPubSubs[proto]
s.PubsubMu.Unlock()
sub, err := topic.Subscribe() // then subscribe to it sub, err := topic.Subscribe() // then subscribe to it
if err != nil { if err != nil {
return err return err
} }
// launch loop waiting for results. // launch loop waiting for results.
go waitResults[T](s, ctx, sub, proto, timeout, f) go waitResults(s, ctx, sub, proto, timeout, f)
return nil return nil
} }
@@ -207,10 +203,5 @@ func waitResults[T interface{}](s *LongLivedPubSubService, ctx context.Context,
continue continue
} }
f(ctx, evt, fmt.Sprintf("%v", proto)) f(ctx, evt, fmt.Sprintf("%v", proto))
/*if p, err := ps.Node.GetPeerRecord(ctx, evt.From); err == nil && len(p) > 0 {
if err := ps.processEvent(ctx, p[0], &evt, topicName); err != nil {
logger.Err(err)
}
}*/
} }
} }

View File

@@ -68,7 +68,7 @@ func (ix *LongLivedStreamRecordedService[T]) gc() {
} }
ix.PubsubMu.Lock() ix.PubsubMu.Lock()
if ix.LongLivedPubSubs[TopicPubSubNodeActivity] != nil { if ix.LongLivedPubSubs[TopicPubSubNodeActivity] != nil {
ad, err := pp.AddrInfoFromString("/ip4/" + conf.GetConfig().Hostname + " /tcp/" + fmt.Sprintf("%v", conf.GetConfig().NodeEndpointPort) + " /p2p/" + ix.Host.ID().String()) ad, err := pp.AddrInfoFromString("/ip4/" + conf.GetConfig().Hostname + "/tcp/" + fmt.Sprintf("%v", conf.GetConfig().NodeEndpointPort) + "/p2p/" + ix.Host.ID().String())
if err == nil { if err == nil {
if b, err := json.Marshal(TopicNodeActivityPub{ if b, err := json.Marshal(TopicNodeActivityPub{
Disposer: *ad, Disposer: *ad,
@@ -235,7 +235,6 @@ var StreamIndexers ProtocolStream = ProtocolStream{}
func ConnectToIndexers(h host.Host, minIndexer int, maxIndexer int, myPID pp.ID) { func ConnectToIndexers(h host.Host, minIndexer int, maxIndexer int, myPID pp.ID) {
logger := oclib.GetLogger() logger := oclib.GetLogger()
ctx := context.Background()
addresses := strings.Split(conf.GetConfig().IndexerAddresses, ",") addresses := strings.Split(conf.GetConfig().IndexerAddresses, ",")
if len(addresses) > maxIndexer { if len(addresses) > maxIndexer {
@@ -243,13 +242,17 @@ func ConnectToIndexers(h host.Host, minIndexer int, maxIndexer int, myPID pp.ID)
} }
for _, indexerAddr := range addresses { for _, indexerAddr := range addresses {
fmt.Println("GENERATE ADDR", indexerAddr)
ad, err := pp.AddrInfoFromString(indexerAddr) ad, err := pp.AddrInfoFromString(indexerAddr)
if err != nil { if err != nil {
fmt.Println("ADDR ERR", err)
logger.Err(err) logger.Err(err)
continue continue
} }
force := false
if h.Network().Connectedness(ad.ID) != network.Connected { if h.Network().Connectedness(ad.ID) != network.Connected {
if err := h.Connect(ctx, *ad); err != nil { force = true
if err := h.Connect(context.Background(), *ad); err != nil {
fmt.Println(err) fmt.Println(err)
logger.Err(err) logger.Err(err)
continue continue
@@ -258,7 +261,7 @@ func ConnectToIndexers(h host.Host, minIndexer int, maxIndexer int, myPID pp.ID)
StaticIndexers = append(StaticIndexers, ad) StaticIndexers = append(StaticIndexers, ad)
// make a privilege streams with indexer. // make a privilege streams with indexer.
for _, proto := range []protocol.ID{ProtocolPublish, ProtocolGet, ProtocolHeartbeat} { for _, proto := range []protocol.ID{ProtocolPublish, ProtocolGet, ProtocolHeartbeat} {
AddStreamProtocol(nil, StreamIndexers, h, proto, ad.ID, myPID, true, nil) AddStreamProtocol(nil, StreamIndexers, h, proto, ad.ID, myPID, force, nil)
} }
} }
if len(StaticIndexers) == 0 { if len(StaticIndexers) == 0 {
@@ -268,7 +271,7 @@ func ConnectToIndexers(h host.Host, minIndexer int, maxIndexer int, myPID pp.ID)
if len(StaticIndexers) < minIndexer { if len(StaticIndexers) < minIndexer {
// TODO : ask for unknown indexer. // TODO : ask for unknown indexer.
} }
SendHeartbeat(ctx, ProtocolHeartbeat, conf.GetConfig().Name, h, StreamIndexers, StaticIndexers, 20*time.Second) // your indexer is just like a node for the next indexer. SendHeartbeat(context.Background(), ProtocolHeartbeat, conf.GetConfig().Name, h, StreamIndexers, StaticIndexers, 20*time.Second) // your indexer is just like a node for the next indexer.
} }
func AddStreamProtocol(ctx *context.Context, protoS ProtocolStream, h host.Host, proto protocol.ID, id pp.ID, mypid pp.ID, force bool, onStreamCreated *func(network.Stream)) ProtocolStream { func AddStreamProtocol(ctx *context.Context, protoS ProtocolStream, h host.Host, proto protocol.ID, id pp.ID, mypid pp.ID, force bool, onStreamCreated *func(network.Stream)) ProtocolStream {
@@ -294,7 +297,7 @@ func AddStreamProtocol(ctx *context.Context, protoS ProtocolStream, h host.Host,
if protoS[proto][id] != nil { if protoS[proto][id] != nil {
protoS[proto][id].Expiry = time.Now().Add(2 * time.Minute) protoS[proto][id].Expiry = time.Now().Add(2 * time.Minute)
} else { } else {
fmt.Println("GENERATE STREAM", proto, id) fmt.Println("NEW STREAM", proto, id)
s, err := h.NewStream(*ctx, id, proto) s, err := h.NewStream(*ctx, id, proto)
if err != nil { if err != nil {
panic(err.Error()) panic(err.Error())
@@ -385,18 +388,3 @@ func sendHeartbeat(ctx context.Context, h host.Host, proto protocol.ID, p *pp.Ad
pss.Expiry = time.Now().UTC().Add(2 * time.Minute) pss.Expiry = time.Now().UTC().Add(2 * time.Minute)
return nil return nil
} }
/*
func SearchPeer(search string) ([]*peer.Peer, error) {
ps := []*peer.Peer{}
access := oclib.NewRequestAdmin(oclib.LibDataEnum(oclib.PEER), nil)
peers := access.Search(nil, search, false)
if len(peers.Data) == 0 {
return ps, errors.New("no self available")
}
for _, p := range peers.Data {
ps = append(ps, p.(*peer.Peer))
}
return ps, nil
}
*/

View File

@@ -79,7 +79,6 @@ func (pr *PeerRecord) ExtractPeer(ourkey string, key string, pubKey crypto.PubKe
if err != nil { if err != nil {
return false, nil, err return false, nil, err
} }
fmt.Println("ExtractPeer MarshalPublicKey")
rel := pp.NONE rel := pp.NONE
if ourkey == key { // at this point is PeerID is same as our... we are... thats our peer INFO if ourkey == key { // at this point is PeerID is same as our... we are... thats our peer INFO
rel = pp.SELF rel = pp.SELF
@@ -106,7 +105,8 @@ func (pr *PeerRecord) ExtractPeer(ourkey string, key string, pubKey crypto.PubKe
if err != nil { if err != nil {
return pp.SELF == p.Relation, nil, err return pp.SELF == p.Relation, nil, err
} }
go tools.NewNATSCaller().SetNATSPub(tools.CREATE_RESOURCE, tools.NATSResponse{ fmt.Println("SENDPEER SELF")
go tools.NewNATSCaller().SetNATSPub(tools.CREATE_PEER, tools.NATSResponse{
FromApp: "oc-discovery", FromApp: "oc-discovery",
Datatype: tools.PEER, Datatype: tools.PEER,
Method: int(tools.CREATE_PEER), Method: int(tools.CREATE_PEER),
@@ -128,6 +128,7 @@ type GetResponse struct {
} }
func (ix *IndexerService) initNodeHandler() { func (ix *IndexerService) initNodeHandler() {
fmt.Println("Node activity")
ix.Host.SetStreamHandler(common.ProtocolHeartbeat, ix.HandleNodeHeartbeat) ix.Host.SetStreamHandler(common.ProtocolHeartbeat, ix.HandleNodeHeartbeat)
ix.Host.SetStreamHandler(common.ProtocolPublish, ix.handleNodePublish) ix.Host.SetStreamHandler(common.ProtocolPublish, ix.handleNodePublish)
ix.Host.SetStreamHandler(common.ProtocolGet, ix.handleNodeGet) ix.Host.SetStreamHandler(common.ProtocolGet, ix.handleNodeGet)
@@ -182,7 +183,7 @@ func (ix *IndexerService) handleNodePublish(s network.Stream) {
} }
if ix.LongLivedPubSubs[common.TopicPubSubNodeActivity] != nil && !rec.NoPub { if ix.LongLivedPubSubs[common.TopicPubSubNodeActivity] != nil && !rec.NoPub {
ad, err := peer.AddrInfoFromString("/ip4/" + conf.GetConfig().Hostname + " /tcp/" + fmt.Sprintf("%v", conf.GetConfig().NodeEndpointPort) + " /p2p/" + ix.Host.ID().String()) ad, err := peer.AddrInfoFromString("/ip4/" + conf.GetConfig().Hostname + "/tcp/" + fmt.Sprintf("%v", conf.GetConfig().NodeEndpointPort) + "/p2p/" + ix.Host.ID().String())
if err == nil { if err == nil {
if b, err := json.Marshal(common.TopicNodeActivityPub{ if b, err := json.Marshal(common.TopicNodeActivityPub{
Disposer: *ad, Disposer: *ad,

View File

@@ -4,12 +4,53 @@ import (
"context" "context"
"encoding/json" "encoding/json"
"fmt" "fmt"
"oc-discovery/daemons/node/common"
oclib "cloud.o-forge.io/core/oc-lib"
"cloud.o-forge.io/core/oc-lib/config"
"cloud.o-forge.io/core/oc-lib/models/peer"
"cloud.o-forge.io/core/oc-lib/tools" "cloud.o-forge.io/core/oc-lib/tools"
pp "github.com/libp2p/go-libp2p/core/peer"
) )
func ListenNATS(n Node) { func ListenNATS(n Node) {
tools.NewNATSCaller().ListenNats(map[tools.NATSMethod]func(tools.NATSResponse){ tools.NewNATSCaller().ListenNats(map[tools.NATSMethod]func(tools.NATSResponse){
tools.CREATE_PEER: func(resp tools.NATSResponse) {
if resp.FromApp == config.GetAppName() {
return
}
logger := oclib.GetLogger()
m := map[string]interface{}{}
err := json.Unmarshal(resp.Payload, &m)
if err != nil {
logger.Err(err)
return
}
p := &peer.Peer{}
p = p.Deserialize(m, p).(*peer.Peer)
ad, err := pp.AddrInfoFromString(p.PeerID)
n.StreamService.Mu.Lock()
defer n.StreamService.Mu.Unlock()
if p.Relation == peer.PARTNER {
n.StreamService.ConnectToPartner(ad.ID, ad)
} else {
ps := common.ProtocolStream{}
for p, s := range n.StreamService.Streams {
m := map[pp.ID]*common.Stream{}
for k := range s {
if ad.ID != k {
m[k] = s[k]
} else {
s[k].Stream.Close()
}
}
ps[p] = m
}
n.StreamService.Streams = ps
}
},
tools.PROPALGATION_EVENT: func(resp tools.NATSResponse) { tools.PROPALGATION_EVENT: func(resp tools.NATSResponse) {
var propalgation tools.PropalgationMessage var propalgation tools.PropalgationMessage
err := json.Unmarshal(resp.Payload, &propalgation) err := json.Unmarshal(resp.Payload, &propalgation)

View File

@@ -237,7 +237,7 @@ func (d *Node) claimInfo(
} }
rec.APIUrl = endPoint rec.APIUrl = endPoint
rec.StreamAddress = "/ip4/" + conf.GetConfig().Hostname + " /tcp/" + fmt.Sprintf("%v", conf.GetConfig().NodeEndpointPort) + " /p2p/" + rec.PeerID rec.StreamAddress = "/ip4/" + conf.GetConfig().Hostname + "/tcp/" + fmt.Sprintf("%v", conf.GetConfig().NodeEndpointPort) + "/p2p/" + rec.PeerID
rec.NATSAddress = oclib.GetConfig().NATSUrl rec.NATSAddress = oclib.GetConfig().NATSUrl
rec.WalletAddress = "my-wallet" rec.WalletAddress = "my-wallet"
rec.ExpiryDate = expiry rec.ExpiryDate = expiry

View File

@@ -91,8 +91,8 @@ func (ps *StreamService) ToPartnerPublishEvent(
if err != nil { if err != nil {
return err return err
} }
ps.mu.Lock() ps.Mu.Lock()
defer ps.mu.Unlock() defer ps.Mu.Unlock()
if p.Relation == peer.PARTNER { if p.Relation == peer.PARTNER {
if ps.Streams[ProtocolHeartbeatPartner] == nil { if ps.Streams[ProtocolHeartbeatPartner] == nil {
ps.Streams[ProtocolHeartbeatPartner] = map[pp.ID]*common.Stream{} ps.Streams[ProtocolHeartbeatPartner] = map[pp.ID]*common.Stream{}
@@ -138,8 +138,8 @@ func (s *StreamService) write(
if dt != nil { if dt != nil {
name = action.String() + "." + (*dt).String() + "#" + peerID.ID.String() name = action.String() + "." + (*dt).String() + "#" + peerID.ID.String()
} }
s.mu.Lock() s.Mu.Lock()
defer s.mu.Unlock() defer s.Mu.Unlock()
if s.Streams[proto] == nil { if s.Streams[proto] == nil {
s.Streams[proto] = map[pp.ID]*common.Stream{} s.Streams[proto] = map[pp.ID]*common.Stream{}
} }

View File

@@ -41,7 +41,7 @@ type StreamService struct {
Node common.DiscoveryPeer Node common.DiscoveryPeer
Streams common.ProtocolStream Streams common.ProtocolStream
maxNodesConn int maxNodesConn int
mu sync.Mutex Mu sync.Mutex
// Stream map[protocol.ID]map[pp.ID]*daemons.Stream // Stream map[protocol.ID]map[pp.ID]*daemons.Stream
} }
@@ -67,8 +67,8 @@ func (s *StreamService) HandlePartnerHeartbeat(stream network.Stream) {
if err != nil { if err != nil {
return return
} }
s.mu.Lock() s.Mu.Lock()
defer s.mu.Unlock() defer s.Mu.Unlock()
if s.Streams[ProtocolHeartbeatPartner] == nil { if s.Streams[ProtocolHeartbeatPartner] == nil {
s.Streams[ProtocolHeartbeatPartner] = map[pp.ID]*common.Stream{} s.Streams[ProtocolHeartbeatPartner] = map[pp.ID]*common.Stream{}
@@ -89,6 +89,20 @@ func (s *StreamService) HandlePartnerHeartbeat(stream network.Stream) {
} }
func (s *StreamService) connectToPartners() error { func (s *StreamService) connectToPartners() error {
for _, proto := range protocols {
f := func(ss network.Stream) {
if s.Streams[proto] == nil {
s.Streams[proto] = map[pp.ID]*common.Stream{}
}
s.Streams[proto][ss.Conn().RemotePeer()] = &common.Stream{
Stream: ss,
Expiry: time.Now().UTC().Add(2 * time.Minute),
}
go s.readLoop(s.Streams[proto][ss.Conn().RemotePeer()])
}
fmt.Println("SetStreamHandler", proto)
s.Host.SetStreamHandler(proto, f)
}
peers, err := s.searchPeer(fmt.Sprintf("%v", peer.PARTNER.EnumIndex())) peers, err := s.searchPeer(fmt.Sprintf("%v", peer.PARTNER.EnumIndex()))
if err != nil { if err != nil {
return err return err
@@ -105,26 +119,13 @@ func (s *StreamService) connectToPartners() error {
s.ConnectToPartner(pid, ad) s.ConnectToPartner(pid, ad)
// heartbeat your partner. // heartbeat your partner.
} }
for _, proto := range protocols {
f := func(ss network.Stream) {
if s.Streams[proto] == nil {
s.Streams[proto] = map[pp.ID]*common.Stream{}
}
s.Streams[proto][ss.Conn().RemotePeer()] = &common.Stream{
Stream: ss,
Expiry: time.Now().UTC().Add(2 * time.Minute),
}
s.readLoop(s.Streams[proto][ss.Conn().RemotePeer()])
}
fmt.Println("SetStreamHandler", proto)
s.Host.SetStreamHandler(proto, f)
}
// TODO if handle... from partner then HeartBeat back // TODO if handle... from partner then HeartBeat back
return nil return nil
} }
func (s *StreamService) ConnectToPartner(pid pp.ID, ad *pp.AddrInfo) { func (s *StreamService) ConnectToPartner(pid pp.ID, ad *pp.AddrInfo) {
logger := oclib.GetLogger() logger := oclib.GetLogger()
force := false
for _, proto := range protocols { for _, proto := range protocols {
f := func(ss network.Stream) { f := func(ss network.Stream) {
if s.Streams[proto] == nil { if s.Streams[proto] == nil {
@@ -134,15 +135,16 @@ func (s *StreamService) ConnectToPartner(pid pp.ID, ad *pp.AddrInfo) {
Stream: ss, Stream: ss,
Expiry: time.Now().UTC().Add(2 * time.Minute), Expiry: time.Now().UTC().Add(2 * time.Minute),
} }
s.readLoop(s.Streams[proto][pid]) go s.readLoop(s.Streams[proto][pid])
} }
if s.Host.Network().Connectedness(ad.ID) != network.Connected { if s.Host.Network().Connectedness(ad.ID) != network.Connected {
force = true
if err := s.Host.Connect(context.Background(), *ad); err != nil { if err := s.Host.Connect(context.Background(), *ad); err != nil {
logger.Err(err) logger.Err(err)
continue continue
} }
} }
s.Streams = common.AddStreamProtocol(nil, s.Streams, s.Host, proto, pid, s.Key, false, &f) s.Streams = common.AddStreamProtocol(nil, s.Streams, s.Host, proto, pid, s.Key, force, &f)
} }
common.SendHeartbeat(context.Background(), ProtocolHeartbeatPartner, conf.GetConfig().Name, common.SendHeartbeat(context.Background(), ProtocolHeartbeatPartner, conf.GetConfig().Name,
s.Host, s.Streams, []*pp.AddrInfo{ad}, 20*time.Second) s.Host, s.Streams, []*pp.AddrInfo{ad}, 20*time.Second)
@@ -153,14 +155,15 @@ func (s *StreamService) searchPeer(search string) ([]*peer.Peer, error) {
ps := []*peer.Peer{} ps := []*peer.Peer{}
if conf.GetConfig().PeerIDS != "" { if conf.GetConfig().PeerIDS != "" {
for _, peerID := range strings.Split(conf.GetConfig().PeerIDS, ",") { for _, peerID := range strings.Split(conf.GetConfig().PeerIDS, ",") {
ppID := strings.Split(peerID, ":") ppID := strings.Split(peerID, "/")
fmt.Println(ppID, peerID)
ps = append(ps, &peer.Peer{ ps = append(ps, &peer.Peer{
AbstractObject: utils.AbstractObject{ AbstractObject: utils.AbstractObject{
UUID: uuid.New().String(), UUID: uuid.New().String(),
Name: ppID[1], Name: ppID[1],
}, },
PeerID: ppID[1], PeerID: ppID[len(ppID)-1],
StreamAddress: "/ip4/127.0.0.1/tcp/" + ppID[0] + "/p2p/" + ppID[1], StreamAddress: peerID,
State: peer.ONLINE, State: peer.ONLINE,
Relation: peer.PARTNER, Relation: peer.PARTNER,
}) })
@@ -194,8 +197,8 @@ func (s *StreamService) StartGC(interval time.Duration) {
} }
func (s *StreamService) gc() { func (s *StreamService) gc() {
s.mu.Lock() s.Mu.Lock()
defer s.mu.Unlock() defer s.Mu.Unlock()
now := time.Now().UTC() now := time.Now().UTC()
if s.Streams[ProtocolHeartbeatPartner] == nil { if s.Streams[ProtocolHeartbeatPartner] == nil {

View File

@@ -4,5 +4,5 @@
"NATS_URL": "nats://nats:4222", "NATS_URL": "nats://nats:4222",
"NODE_MODE": "indexer", "NODE_MODE": "indexer",
"NODE_ENDPOINT_PORT": 4002, "NODE_ENDPOINT_PORT": 4002,
"INDEXER_ADDRESSES": "/ip4/oc-discovery1/tcp/4001/p2p/12D3KooWGn3j4XqTSrjJDGGpTQERdDV5TPZdhQp87rAUnvQssvQu" "INDEXER_ADDRESSES": "/ip4/172.19.0.2/tcp/4001/p2p/12D3KooWGn3j4XqTSrjJDGGpTQERdDV5TPZdhQp87rAUnvQssvQu"
} }

View File

@@ -4,5 +4,5 @@
"NATS_URL": "nats://nats:4222", "NATS_URL": "nats://nats:4222",
"NODE_MODE": "node", "NODE_MODE": "node",
"NODE_ENDPOINT_PORT": 4003, "NODE_ENDPOINT_PORT": 4003,
"INDEXER_ADDRESSES": "/ip4/oc-discovery2/tcp/4002/p2p/12D3KooWC3GNStak8KCYtJq11Dxiq45EJV53z1ZvKetMcZBeBX6u" "INDEXER_ADDRESSES": "/ip4/172.19.0.3/tcp/4002/p2p/12D3KooWC3GNStak8KCYtJq11Dxiq45EJV53z1ZvKetMcZBeBX6u"
} }

View File

@@ -4,6 +4,6 @@
"NATS_URL": "nats://nats:4222", "NATS_URL": "nats://nats:4222",
"NODE_MODE": "node", "NODE_MODE": "node",
"NODE_ENDPOINT_PORT": 4004, "NODE_ENDPOINT_PORT": 4004,
"INDEXER_ADDRESSES": "/ip4/oc-discovery1/tcp/4001/p2p/12D3KooWGn3j4XqTSrjJDGGpTQERdDV5TPZdhQp87rAUnvQssvQu", "INDEXER_ADDRESSES": "/ip4/172.19.0.2/tcp/4001/p2p/12D3KooWGn3j4XqTSrjJDGGpTQERdDV5TPZdhQp87rAUnvQssvQu",
"PEER_IDS": "/ip4/oc-discovery3/tcp/4003/p2p/12D3KooWGn3j4XqTSrjJDGGpTQERdDV5TPZdhQp87rAUnvQssvQu" "PEER_IDS": "/ip4/172.19.0.4/tcp/4003/p2p/12D3KooWBh9kZrekBAE5G33q4jCLNRAzygem3gP1mMdK8mhoCTaw"
} }

2
go.mod
View File

@@ -3,7 +3,7 @@ module oc-discovery
go 1.24.6 go 1.24.6
require ( require (
cloud.o-forge.io/core/oc-lib v0.0.0-20260203150531-ef916fe2d995 cloud.o-forge.io/core/oc-lib v0.0.0-20260205131630-342451db2581
github.com/beego/beego v1.12.13 github.com/beego/beego v1.12.13
github.com/beego/beego/v2 v2.3.8 github.com/beego/beego/v2 v2.3.8
github.com/go-redis/redis v6.15.9+incompatible github.com/go-redis/redis v6.15.9+incompatible

4
go.sum
View File

@@ -34,6 +34,10 @@ cloud.o-forge.io/core/oc-lib v0.0.0-20260203150123-4258f6b58083 h1:nKiU4AfeX+axS
cloud.o-forge.io/core/oc-lib v0.0.0-20260203150123-4258f6b58083/go.mod h1:T0UCxRd8w+qCVVC0NEyDiWIGC5ADwEbQ7hFcvftd4Ks= cloud.o-forge.io/core/oc-lib v0.0.0-20260203150123-4258f6b58083/go.mod h1:T0UCxRd8w+qCVVC0NEyDiWIGC5ADwEbQ7hFcvftd4Ks=
cloud.o-forge.io/core/oc-lib v0.0.0-20260203150531-ef916fe2d995 h1:ZDRvnzTTNHgMm5hYmseHdEPqQ6rn/4v+P9f/JIxPaNw= cloud.o-forge.io/core/oc-lib v0.0.0-20260203150531-ef916fe2d995 h1:ZDRvnzTTNHgMm5hYmseHdEPqQ6rn/4v+P9f/JIxPaNw=
cloud.o-forge.io/core/oc-lib v0.0.0-20260203150531-ef916fe2d995/go.mod h1:T0UCxRd8w+qCVVC0NEyDiWIGC5ADwEbQ7hFcvftd4Ks= cloud.o-forge.io/core/oc-lib v0.0.0-20260203150531-ef916fe2d995/go.mod h1:T0UCxRd8w+qCVVC0NEyDiWIGC5ADwEbQ7hFcvftd4Ks=
cloud.o-forge.io/core/oc-lib v0.0.0-20260205131048-425cd2a9ba2f h1:Ku6u+SeoNXHMBzckekGyXCHLDJPh20Y8GayO6fXEcZE=
cloud.o-forge.io/core/oc-lib v0.0.0-20260205131048-425cd2a9ba2f/go.mod h1:T0UCxRd8w+qCVVC0NEyDiWIGC5ADwEbQ7hFcvftd4Ks=
cloud.o-forge.io/core/oc-lib v0.0.0-20260205131630-342451db2581 h1:V9eANWFEkoEPg3nWCvYXnLYbKDdAm3/Y7uCw1nt22Cc=
cloud.o-forge.io/core/oc-lib v0.0.0-20260205131630-342451db2581/go.mod h1:T0UCxRd8w+qCVVC0NEyDiWIGC5ADwEbQ7hFcvftd4Ks=
github.com/BurntSushi/toml v0.3.1/go.mod h1:xHWCNGjB5oqiDr8zfno3MHue2Ht5sIBksp03qcyfWMU= github.com/BurntSushi/toml v0.3.1/go.mod h1:xHWCNGjB5oqiDr8zfno3MHue2Ht5sIBksp03qcyfWMU=
github.com/Knetic/govaluate v3.0.0+incompatible/go.mod h1:r7JcOSlj0wfOMncg0iLm8Leh48TZaKVeNIfJntJ2wa0= github.com/Knetic/govaluate v3.0.0+incompatible/go.mod h1:r7JcOSlj0wfOMncg0iLm8Leh48TZaKVeNIfJntJ2wa0=
github.com/alecthomas/template v0.0.0-20160405071501-a0175ee3bccc/go.mod h1:LOuyumcjzFXgccqObfd/Ljyb9UuFJ6TxHnclSeseNhc= github.com/alecthomas/template v0.0.0-20160405071501-a0175ee3bccc/go.mod h1:LOuyumcjzFXgccqObfd/Ljyb9UuFJ6TxHnclSeseNhc=

View File

@@ -21,7 +21,7 @@ func main() {
oclib.InitDaemon(appname) oclib.InitDaemon(appname)
// get the right config file // get the right config file
o := oclib.GetConfLoader() o := oclib.GetConfLoader(appname)
conf.GetConfig().Name = o.GetStringDefault("NAME", "opencloud-demo") conf.GetConfig().Name = o.GetStringDefault("NAME", "opencloud-demo")
conf.GetConfig().Hostname = o.GetStringDefault("HOSTNAME", "127.0.0.1") conf.GetConfig().Hostname = o.GetStringDefault("HOSTNAME", "127.0.0.1")