blob: 791ce0a49c8a55f34d8e6921c938a0cb10bc9809 [file]
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package metadata
import (
"context"
"encoding/json"
"errors"
"fmt"
"reflect"
)
import (
"github.com/dubbogo/gost/log/logger"
)
import (
"dubbo.apache.org/dubbo-go/v3/common"
"dubbo.apache.org/dubbo-go/v3/common/constant"
"dubbo.apache.org/dubbo-go/v3/common/extension"
"dubbo.apache.org/dubbo-go/v3/metadata/info"
tripleapi "dubbo.apache.org/dubbo-go/v3/metadata/triple_api/proto"
"dubbo.apache.org/dubbo-go/v3/protocol/base"
"dubbo.apache.org/dubbo-go/v3/protocol/invocation"
"dubbo.apache.org/dubbo-go/v3/registry"
)
const defaultTimeout = "5s" // s
func GetMetadataFromMetadataReport(revision string, instance registry.ServiceInstance, registryId string) (*info.MetadataInfo, error) {
report := GetMetadataReportByRegistry(registryId)
if report == nil {
return nil, fmt.Errorf("metadata_report failed: operation=get app=%s revision=%s registry_id=%s storage_type=%s: no metadata report instance found, please check metadata-report configuration",
instance.GetServiceName(), revision, registryId, constant.RemoteMetadataStorageType)
}
meta, err := report.GetAppMetadata(instance.GetServiceName(), revision)
if err != nil {
return nil, fmt.Errorf("%w; registry_id=%s", err, registryId)
}
return meta, nil
}
func GetMetadataFromRpc(revision string, instance registry.ServiceInstance) (*info.MetadataInfo, error) {
return GetMetadataFromRpcWithContext(context.Background(), revision, instance)
}
// GetMetadataFromRpcWithContext fetches metadata through the metadata service
// while preserving the caller's context for the underlying RPC invocation.
func GetMetadataFromRpcWithContext(ctx context.Context, revision string, instance registry.ServiceInstance) (*info.MetadataInfo, error) {
if ctx == nil {
ctx = context.Background()
}
storageType := constant.DefaultMetadataStorageType
if instanceMetadata := instance.GetMetadata(); instanceMetadata != nil && instanceMetadata[constant.MetadataStorageTypePropertyName] != "" {
storageType = instanceMetadata[constant.MetadataStorageTypePropertyName]
}
if err := ctx.Err(); err != nil {
return nil, fmt.Errorf("rpc_metadata failed: app=%s revision=%s instance_id=%s host=%s storage_type=%s: %w",
instance.GetServiceName(), revision, instance.GetID(), instance.GetHost(), storageType, err)
}
url, err := buildStandardMetadataServiceURL(instance)
if err != nil {
return nil, fmt.Errorf("url_construction failed: app=%s revision=%s instance_id=%s host=%s storage_type=%s: %w",
instance.GetServiceName(), revision, instance.GetID(), instance.GetHost(), storageType, err)
}
url.SetParam(constant.TimeoutKey, defaultTimeout)
p := extension.GetProtocol(url.Protocol)
invoker := p.Refer(url)
if invoker == nil { // can't connect instance
return nil, fmt.Errorf("rpc_metadata failed: app=%s revision=%s instance_id=%s host=%s storage_type=%s: can not connect to remote metadata service",
instance.GetServiceName(), revision, instance.GetID(), instance.GetHost(), storageType)
}
var remoteService remoteMetadataService
if url.Protocol == constant.TriProtocol && instance.GetMetadata()[constant.MetadataVersion] == constant.MetadataServiceV2Version {
remoteService = &triMetadataServiceV2{invoker: invoker}
} else {
remoteService = &remoteMetadataServiceV1{invoker: invoker}
}
defer func() {
invoker.Destroy()
}()
metadataInfo, err := remoteService.getMetadataInfo(ctx, revision)
if err != nil {
return metadataInfo, fmt.Errorf("rpc_metadata failed: app=%s revision=%s instance_id=%s host=%s storage_type=%s: %w",
instance.GetServiceName(), revision, instance.GetID(), instance.GetHost(), storageType, err)
}
return metadataInfo, nil
}
// remoteMetadataService is the internal interface for fetching MetadataInfo via RPC.
type remoteMetadataService interface {
getMetadataInfo(ctx context.Context, revision string) (*info.MetadataInfo, error)
}
type triMetadataServiceV2 struct {
invoker base.Invoker
}
// getMetadataInfo fetches metadata via RPC using the Triple protocol (Protobuf).
func (m *triMetadataServiceV2) getMetadataInfo(ctx context.Context, revision string) (*info.MetadataInfo, error) {
const methodName = "GetMetadataInfo"
req := &tripleapi.MetadataRequest{Revision: revision}
metadataInfo := &tripleapi.MetadataInfoV2{}
inv, _ := generateInvocation(m.invoker.GetURL(), methodName, req, metadataInfo, constant.CallUnary)
if rpcInv, ok := inv.(*invocation.RPCInvocation); ok {
rpcInv.SetContext(ctx)
}
res := m.invoker.Invoke(ctx, inv)
if res.Error() != nil {
logger.Errorf("[Metadata][RPC] could not get the metadata info from remote provider, err=%v", res.Error())
return nil, fmt.Errorf("remote metadata call failed: %w", res.Error())
}
return convertMetadataInfoV2(metadataInfo), nil
}
func convertMetadataInfoV2(v2 *tripleapi.MetadataInfoV2) *info.MetadataInfo {
infos := make(map[string]*info.ServiceInfo, 0)
for k, v := range v2.Services {
serviceInfo := &info.ServiceInfo{
Name: v.Name,
Group: v.Group,
Version: v.Version,
Protocol: v.Protocol,
Path: v.Path,
Params: v.Params,
}
infos[k] = serviceInfo
}
metadataInfo := &info.MetadataInfo{
App: v2.App,
Revision: v2.Version,
Tag: v2.Tag,
Services: infos,
}
return metadataInfo
}
func generateInvocation(u *common.URL, methodName string, req any, resp any, callType string) (base.Invocation, error) {
var inv *invocation.RPCInvocation
if u.Protocol == constant.TriProtocol {
var paramsRawVals []any
paramsRawVals = append(paramsRawVals, req)
if resp != nil {
paramsRawVals = append(paramsRawVals, resp)
}
inv = invocation.NewRPCInvocationWithOptions(
invocation.WithMethodName(methodName),
invocation.WithAttachment(constant.TimeoutKey, "5000"),
invocation.WithAttachment(constant.RetriesKey, "2"),
invocation.WithArguments([]any{req}),
invocation.WithReply(resp),
invocation.WithParameterRawValues(paramsRawVals),
)
inv.SetAttribute(constant.CallTypeKey, callType)
} else {
rV := reflect.ValueOf(req)
inv = invocation.NewRPCInvocationWithOptions(
invocation.WithMethodName(methodName),
invocation.WithArguments([]any{rV.Interface()}),
invocation.WithReply(resp),
invocation.WithAttachments(map[string]any{constant.AsyncKey: "false"}),
invocation.WithParameterValues([]reflect.Value{rV}))
}
return inv, nil
}
type remoteMetadataServiceV1 struct {
invoker base.Invoker
}
// getMetadataInfo fetches metadata via RPC using the dubbo:// protocol (Hessian2 serialization).
func (m *remoteMetadataServiceV1) getMetadataInfo(ctx context.Context, revision string) (*info.MetadataInfo, error) {
const methodName = "getMetadataInfo"
// Use interface{} as reply parameter to accept any type (MetadataInfo or string)
// This avoids panic when Java returns String instead of MetadataInfo
var rawResult any
inv, _ := generateInvocation(m.invoker.GetURL(), methodName, revision, &rawResult, constant.CallUnary)
if rpcInv, ok := inv.(*invocation.RPCInvocation); ok {
rpcInv.SetContext(ctx)
}
res := m.invoker.Invoke(ctx, inv)
if res.Error() != nil {
logger.Errorf("[Metadata][RPC] RPC call failed to %s, err=%v", m.invoker.GetURL().Location, res.Error())
return nil, fmt.Errorf("RPC call failed to %s: %w", m.invoker.GetURL().Location, res.Error())
}
// rawResult now contains the deserialized value - could be *MetadataInfo, string, or nil
// Handle nil response (e.g., Java service not fully initialized)
if rawResult == nil {
logger.Warnf("[Metadata][RPC] Provider %s returned nil metadata (service may not be ready), revision=%s",
m.invoker.GetURL().Location, revision)
return nil, fmt.Errorf("metadata is nil from %s, revision: %s", m.invoker.GetURL().Location, revision)
}
var metadataInfo *info.MetadataInfo
// Try to handle different return types from Java Dubbo
if result, ok := rawResult.(*info.MetadataInfo); ok {
metadataInfo = result
} else if strValue, ok := rawResult.(string); ok {
// Old Java Dubbo version returns JSON string instead of MetadataInfo object
// Try to parse it as JSON for backward compatibility
logger.Warnf("[Metadata][RPC] Provider %s returned string type (old Dubbo version), attempting JSON parse", m.invoker.GetURL().Location)
metadataInfo = &info.MetadataInfo{}
if err := json.Unmarshal([]byte(strValue), metadataInfo); err != nil {
logger.Errorf("[Metadata][RPC] failed to parse JSON string from provider %s, err=%v", m.invoker.GetURL().Location, err)
logger.Errorf("[Metadata][RPC] - String content: %s", truncateString(strValue, 1000))
return nil, fmt.Errorf("failed to parse metadata JSON from %s: %v", m.invoker.GetURL().Location, err)
}
} else {
// Neither MetadataInfo nor String - this is unexpected
logger.Errorf("[Metadata][RPC] unexpected metadata type from %s: got %T, expected *info.MetadataInfo or string",
m.invoker.GetURL().Location, rawResult)
return nil, fmt.Errorf("unexpected metadata type from %s: got %T, expected *info.MetadataInfo or string",
m.invoker.GetURL().Location, rawResult)
}
return metadataInfo, nil
}
// truncateString truncates a string to maxLen characters
func truncateString(s string, maxLen int) string {
if len(s) <= maxLen {
return s
}
return s[:maxLen] + "..."
}
// buildStandardMetadataServiceURL will use standard format to build the metadata service url.
// Returns an error if required params (protocol or port) are missing.
func buildStandardMetadataServiceURL(ins registry.ServiceInstance) (*common.URL, error) {
ps := getMetadataServiceUrlParams(ins)
if ps[constant.ProtocolKey] == "" {
return nil, errors.New("metadata service URL params missing: protocol is empty")
}
if ps[constant.PortKey] == "" {
return nil, errors.New("metadata service URL params missing: port is empty")
}
sn := ins.GetServiceName()
host := ins.GetHost()
metaV := ins.GetMetadata()[constant.MetadataVersion]
proto := ps[constant.ProtocolKey]
convertedParams := make(map[string][]string, len(ps))
for k, v := range ps {
convertedParams[k] = []string{v}
}
u := common.NewURLWithOptions(common.WithIp(host),
common.WithPath(constant.MetadataServiceName),
common.WithProtocol(proto),
common.WithPort(ps[constant.PortKey]),
common.WithParams(convertedParams),
common.WithParamsValue(constant.GroupKey, sn),
common.WithParamsValue(constant.InterfaceKey, constant.MetadataServiceName))
if proto == constant.TriProtocol {
u.SetAttribute(constant.ClientInfoKey, "info")
u.Methods = []string{"GetMetadataInfo", "getMetadataInfo"}
if metaV == constant.MetadataServiceV2Version {
u.Path = constant.MetadataServiceV2Name
u.SetParam(constant.VersionKey, metaV)
u.SetParam(constant.InterfaceKey, constant.MetadataServiceV2Name)
u.DelParam(constant.SerializationKey)
} else {
u.SetParam(constant.SerializationKey, constant.Hessian2Serialization)
}
}
return u, nil
}
// getMetadataServiceUrlParams this will convertV2 the metadata service url parameters to map structure
// it looks like:
// {"dubbo":{"timeout":"10000","version":"1.0.0","dubbo":"2.0.2","release":"2.7.6","port":"20880"}}
func getMetadataServiceUrlParams(ins registry.ServiceInstance) map[string]string {
ps := ins.GetMetadata()
res := make(map[string]string, 2)
if str, ok := ps[constant.MetadataServiceURLParamsPropertyName]; ok && len(str) > 0 {
err := json.Unmarshal([]byte(str), &res)
if err != nil {
logger.Errorf("[Metadata][URL] url_construction failed: app=%s instance_id=%s host=%s: could not parse metadata service URL parameters: %v",
ins.GetServiceName(), ins.GetID(), ins.GetHost(), err)
}
}
return res
}