| /* |
| * 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 |
| } |