502 lines
19 KiB
Go
502 lines
19 KiB
Go
package oci
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/oracle/oci-go-sdk/v65/common"
|
|
"github.com/oracle/oci-go-sdk/v65/identity"
|
|
"github.com/oracle/oci-go-sdk/v65/ons"
|
|
"github.com/oracle/oci-go-sdk/v65/sch"
|
|
)
|
|
|
|
// 日志回传链路(方案A)在租户侧的资源命名与描述由 logrelay_names.go 集中派生;
|
|
// 资源建在凭据默认区域,Audit 日志组含子 compartment。
|
|
// Topic 删除后名称有保留期(同名重建长时间 409 Conflict),故 Topic 用
|
|
// 前缀+随机后缀命名、按前缀幂等查找,销毁重建不受保留期阻塞。
|
|
const (
|
|
relayAuditLogGroup = "_Audit_Include_Subcompartment"
|
|
relayTopicPages = 5 // 按前缀查找 Topic 的翻页上限
|
|
)
|
|
|
|
// Connector 创建为异步 work request,轮询直至 ACTIVE。
|
|
var (
|
|
relayConnectorPollTick = 5 * time.Second
|
|
relayConnectorPollLimit = 18
|
|
)
|
|
|
|
// RelayResource 是链路单个云资源的状态;Created 标记本次调用新建,供失败回滚判据。
|
|
type RelayResource struct {
|
|
ID string `json:"id"`
|
|
State string `json:"state"`
|
|
Created bool `json:"-"`
|
|
}
|
|
|
|
// RelayState 是回传链路四资源现状;已删除的资源不出现(ID 为空)。
|
|
type RelayState struct {
|
|
Topic RelayResource `json:"topic"`
|
|
Subscription RelayResource `json:"subscription"`
|
|
Connector RelayResource `json:"connector"`
|
|
Policy RelayResource `json:"policy"`
|
|
}
|
|
|
|
// RelayEventCondition 生成 Connector Log Filter 条件:按 eventName 白名单收窄回传。
|
|
func RelayEventCondition(events []string) string {
|
|
terms := make([]string, 0, len(events))
|
|
for _, e := range events {
|
|
terms = append(terms, fmt.Sprintf("data.eventName='%s'", e))
|
|
}
|
|
return strings.Join(terms, " or ")
|
|
}
|
|
|
|
func (c *RealClient) onsControlClient(cred Credentials) (ons.NotificationControlPlaneClient, error) {
|
|
cp, err := ons.NewNotificationControlPlaneClientWithConfigurationProvider(provider(cred))
|
|
if err != nil {
|
|
return cp, fmt.Errorf("new ons control client: %w", err)
|
|
}
|
|
applyProxy(&cp.BaseClient, cred)
|
|
return cp, nil
|
|
}
|
|
|
|
func (c *RealClient) onsDataClient(cred Credentials) (ons.NotificationDataPlaneClient, error) {
|
|
dp, err := ons.NewNotificationDataPlaneClientWithConfigurationProvider(provider(cred))
|
|
if err != nil {
|
|
return dp, fmt.Errorf("new ons data client: %w", err)
|
|
}
|
|
applyProxy(&dp.BaseClient, cred)
|
|
return dp, nil
|
|
}
|
|
|
|
func (c *RealClient) schClient(cred Credentials) (sch.ServiceConnectorClient, error) {
|
|
sc, err := sch.NewServiceConnectorClientWithConfigurationProvider(provider(cred))
|
|
if err != nil {
|
|
return sc, fmt.Errorf("new sch client: %w", err)
|
|
}
|
|
applyProxy(&sc.BaseClient, cred)
|
|
return sc, nil
|
|
}
|
|
|
|
// EnsureRelayTopic 实现 Client:按新命名前缀返回既有 Topic,fallback 到 legacy 前缀;
|
|
// 未命中则以新前缀 + 随机后缀新建。命中旧命名时顺手把描述刷新为中性文案。
|
|
func (c *RealClient) EnsureRelayTopic(ctx context.Context, cred Credentials) (RelayResource, error) {
|
|
cp, err := c.onsControlClient(cred)
|
|
if err != nil {
|
|
return RelayResource{}, err
|
|
}
|
|
names := relayResourceNames(cred.TenancyOCID)
|
|
res, ok, err := findRelayTopic(ctx, cp, cred.TenancyOCID, names.TopicPrefix, legacyRelayTopicPrefix)
|
|
if err != nil {
|
|
return RelayResource{}, err
|
|
}
|
|
if ok {
|
|
refreshRelayTopicDesc(ctx, cp, res.ID)
|
|
return res, nil
|
|
}
|
|
return createRelayTopic(ctx, cp, cred.TenancyOCID, names.TopicPrefix)
|
|
}
|
|
|
|
// createRelayTopic 以指定前缀 + 4 字节随机后缀新建 ONS Topic,描述使用中性文案。
|
|
func createRelayTopic(ctx context.Context, cp ons.NotificationControlPlaneClient, tenancy, prefix string) (RelayResource, error) {
|
|
name, err := relayTopicNewName(prefix)
|
|
if err != nil {
|
|
return RelayResource{}, err
|
|
}
|
|
created, err := cp.CreateTopic(ctx, ons.CreateTopicRequest{
|
|
CreateTopicDetails: ons.CreateTopicDetails{
|
|
Name: &name,
|
|
CompartmentId: &tenancy,
|
|
Description: common.String(relayTopicDescNew),
|
|
},
|
|
})
|
|
if err != nil {
|
|
return RelayResource{}, fmt.Errorf("create ons topic: %w", err)
|
|
}
|
|
return RelayResource{ID: deref(created.TopicId), State: string(created.LifecycleState), Created: true}, nil
|
|
}
|
|
|
|
// findRelayTopic 依次按传入的前缀列表在租户范围内查找存活 Topic;第一个命中即返回。
|
|
// 支持传入新命名前缀与 legacy 前缀,实现向后兼容存量资源。
|
|
func findRelayTopic(ctx context.Context, cp ons.NotificationControlPlaneClient, tenancy string, prefixes ...string) (RelayResource, bool, error) {
|
|
req := ons.ListTopicsRequest{CompartmentId: &tenancy}
|
|
for page := 0; page < relayTopicPages; page++ {
|
|
list, err := cp.ListTopics(ctx, req)
|
|
if err != nil {
|
|
return RelayResource{}, false, fmt.Errorf("list ons topics: %w", err)
|
|
}
|
|
if res, ok := matchRelayTopic(list.Items, prefixes); ok {
|
|
return res, true, nil
|
|
}
|
|
if list.OpcNextPage == nil {
|
|
break
|
|
}
|
|
req.Page = list.OpcNextPage
|
|
}
|
|
return RelayResource{}, false, nil
|
|
}
|
|
|
|
// matchRelayTopic 在一页 Topic 中挑出第一个命名前缀命中且处于 ACTIVE 的资源。
|
|
func matchRelayTopic(items []ons.NotificationTopicSummary, prefixes []string) (RelayResource, bool) {
|
|
for _, t := range items {
|
|
if t.LifecycleState != ons.NotificationTopicSummaryLifecycleStateActive {
|
|
continue
|
|
}
|
|
name := deref(t.Name)
|
|
for _, p := range prefixes {
|
|
if strings.HasPrefix(name, p) {
|
|
return RelayResource{ID: deref(t.TopicId), State: string(t.LifecycleState)}, true
|
|
}
|
|
}
|
|
}
|
|
return RelayResource{}, false
|
|
}
|
|
|
|
// refreshRelayTopicDesc 尽力把 Topic 描述改为中性文案;失败不阻塞主流程(权限不足等场景直接忽略)。
|
|
func refreshRelayTopicDesc(ctx context.Context, cp ons.NotificationControlPlaneClient, topicID string) {
|
|
if topicID == "" {
|
|
return
|
|
}
|
|
_, _ = cp.UpdateTopic(ctx, ons.UpdateTopicRequest{
|
|
TopicId: &topicID,
|
|
TopicAttributesDetails: ons.TopicAttributesDetails{
|
|
Description: common.String(relayTopicDescNew),
|
|
},
|
|
})
|
|
}
|
|
|
|
// EnsureRelaySubscription 实现 Client:按 endpoint 返回既有 CUSTOM_HTTPS 订阅或新建;
|
|
// 新建订阅处于 PENDING,待 ONS 向 endpoint 投递确认消息、面板回访后转 ACTIVE。
|
|
func (c *RealClient) EnsureRelaySubscription(ctx context.Context, cred Credentials, topicID, endpoint string) (RelayResource, error) {
|
|
dp, err := c.onsDataClient(cred)
|
|
if err != nil {
|
|
return RelayResource{}, err
|
|
}
|
|
if res, ok, err := findRelaySubscription(ctx, dp, cred.TenancyOCID, topicID, endpoint); err != nil || ok {
|
|
return res, err
|
|
}
|
|
created, err := dp.CreateSubscription(ctx, ons.CreateSubscriptionRequest{
|
|
CreateSubscriptionDetails: ons.CreateSubscriptionDetails{
|
|
TopicId: &topicID,
|
|
CompartmentId: &cred.TenancyOCID,
|
|
Protocol: common.String("CUSTOM_HTTPS"),
|
|
Endpoint: &endpoint,
|
|
},
|
|
})
|
|
if err != nil {
|
|
return RelayResource{}, fmt.Errorf("create ons subscription: %w", err)
|
|
}
|
|
return RelayResource{ID: deref(created.Id), State: string(created.LifecycleState), Created: true}, nil
|
|
}
|
|
|
|
// findRelaySubscription 在 Topic 下按 endpoint 查找存活订阅。
|
|
func findRelaySubscription(ctx context.Context, dp ons.NotificationDataPlaneClient, tenancy, topicID, endpoint string) (RelayResource, bool, error) {
|
|
if endpoint == "" {
|
|
return RelayResource{}, false, nil
|
|
}
|
|
list, err := dp.ListSubscriptions(ctx, ons.ListSubscriptionsRequest{
|
|
CompartmentId: &tenancy, TopicId: &topicID,
|
|
})
|
|
if err != nil {
|
|
return RelayResource{}, false, fmt.Errorf("list ons subscriptions: %w", err)
|
|
}
|
|
for _, s := range list.Items {
|
|
if deref(s.Endpoint) == endpoint && s.LifecycleState != ons.SubscriptionSummaryLifecycleStateDeleted {
|
|
return RelayResource{ID: deref(s.Id), State: string(s.LifecycleState)}, true, nil
|
|
}
|
|
}
|
|
return RelayResource{}, false, nil
|
|
}
|
|
|
|
// GetRelaySubscription 实现 Client:查询订阅当前状态,供确认期轮询。
|
|
func (c *RealClient) GetRelaySubscription(ctx context.Context, cred Credentials, subscriptionID string) (RelayResource, error) {
|
|
dp, err := c.onsDataClient(cred)
|
|
if err != nil {
|
|
return RelayResource{}, err
|
|
}
|
|
resp, err := dp.GetSubscription(ctx, ons.GetSubscriptionRequest{SubscriptionId: &subscriptionID})
|
|
if err != nil {
|
|
return RelayResource{}, fmt.Errorf("get ons subscription: %w", err)
|
|
}
|
|
return RelayResource{ID: deref(resp.Id), State: string(resp.LifecycleState)}, nil
|
|
}
|
|
|
|
// EnsureRelayPolicy 实现 Client:按新命名返回既有 IAM policy,fallback 到 legacy 命名;
|
|
// 未命中则以新命名创建。写操作必须发往 home region,授权 Service Connector 发布消息到 Topic。
|
|
func (c *RealClient) EnsureRelayPolicy(ctx context.Context, cred Credentials, homeRegion string) (RelayResource, error) {
|
|
ic, err := c.identityClientAt(cred, homeRegion)
|
|
if err != nil {
|
|
return RelayResource{}, err
|
|
}
|
|
names := relayResourceNames(cred.TenancyOCID)
|
|
res, ok, err := findRelayPolicy(ctx, ic, cred.TenancyOCID, names.PolicyName, legacyRelayPolicyName)
|
|
if err != nil {
|
|
return RelayResource{}, err
|
|
}
|
|
if ok {
|
|
refreshRelayPolicyDesc(ctx, ic, res.ID)
|
|
return res, nil
|
|
}
|
|
return createRelayPolicy(ctx, ic, cred.TenancyOCID, names.PolicyName)
|
|
}
|
|
|
|
// findRelayPolicy 依次按传入的名称列表在租户范围内查找 IAM Policy;第一个命中即返回。
|
|
func findRelayPolicy(ctx context.Context, ic identity.IdentityClient, tenancy string, names ...string) (RelayResource, bool, error) {
|
|
for _, name := range names {
|
|
list, err := ic.ListPolicies(ctx, identity.ListPoliciesRequest{
|
|
CompartmentId: &tenancy, Name: common.String(name),
|
|
})
|
|
if err != nil {
|
|
return RelayResource{}, false, fmt.Errorf("list policies: %w", err)
|
|
}
|
|
if len(list.Items) > 0 {
|
|
it := list.Items[0]
|
|
return RelayResource{ID: deref(it.Id), State: string(it.LifecycleState)}, true, nil
|
|
}
|
|
}
|
|
return RelayResource{}, false, nil
|
|
}
|
|
|
|
// createRelayPolicy 以中性描述与派生名称新建 Policy,statement 允许 Service Connector 发布消息到租户内 Topic。
|
|
func createRelayPolicy(ctx context.Context, ic identity.IdentityClient, tenancy, name string) (RelayResource, error) {
|
|
stmt := fmt.Sprintf("Allow any-user to use ons-topics in tenancy where all {request.principal.type='serviceconnector', request.principal.compartment.id='%s'}", tenancy)
|
|
created, err := ic.CreatePolicy(ctx, identity.CreatePolicyRequest{
|
|
CreatePolicyDetails: identity.CreatePolicyDetails{
|
|
CompartmentId: &tenancy,
|
|
Name: common.String(name),
|
|
Description: common.String(relayPolicyDescNew),
|
|
Statements: []string{stmt},
|
|
},
|
|
})
|
|
if err != nil {
|
|
return RelayResource{}, fmt.Errorf("create policy: %w", err)
|
|
}
|
|
return RelayResource{ID: deref(created.Id), State: string(created.LifecycleState), Created: true}, nil
|
|
}
|
|
|
|
// refreshRelayPolicyDesc 尽力把 Policy 描述改为中性文案;失败不阻塞主流程。
|
|
func refreshRelayPolicyDesc(ctx context.Context, ic identity.IdentityClient, id string) {
|
|
if id == "" {
|
|
return
|
|
}
|
|
_, _ = ic.UpdatePolicy(ctx, identity.UpdatePolicyRequest{
|
|
PolicyId: &id,
|
|
UpdatePolicyDetails: identity.UpdatePolicyDetails{
|
|
Description: common.String(relayPolicyDescNew),
|
|
},
|
|
})
|
|
}
|
|
|
|
// EnsureRelayConnector 实现 Client:按新命名返回既有 Connector,fallback 到 legacy DisplayName;
|
|
// 命中时对齐 Log Filter 条件,未命中则以新命名新建并轮询至 ACTIVE。
|
|
func (c *RealClient) EnsureRelayConnector(ctx context.Context, cred Credentials, topicID, condition string) (RelayResource, error) {
|
|
sc, err := c.schClient(cred)
|
|
if err != nil {
|
|
return RelayResource{}, err
|
|
}
|
|
names := relayResourceNames(cred.TenancyOCID)
|
|
res, ok, err := findRelayConnector(ctx, sc, cred.TenancyOCID, names.ConnectorName, legacyRelayConnectorName)
|
|
if err != nil {
|
|
return RelayResource{}, err
|
|
}
|
|
if ok {
|
|
return res, reconcileRelayCondition(ctx, sc, res.ID, condition)
|
|
}
|
|
return createRelayConnector(ctx, sc, cred.TenancyOCID, names.ConnectorName, topicID, condition)
|
|
}
|
|
|
|
// createRelayConnector 以派生名称新建 Service Connector(_Audit 含子区间 → Topic,按 condition 过滤),
|
|
// 新建后轮询至 ACTIVE,超时返回错误但保留 Created 供上层回滚。Connector 不设 Description,避免恒定文案指纹。
|
|
func createRelayConnector(ctx context.Context, sc sch.ServiceConnectorClient, tenancy, name, topicID, condition string) (RelayResource, error) {
|
|
details := sch.CreateServiceConnectorDetails{
|
|
DisplayName: common.String(name),
|
|
CompartmentId: &tenancy,
|
|
Source: sch.LoggingSourceDetails{LogSources: []sch.LogSource{{
|
|
CompartmentId: &tenancy,
|
|
LogGroupId: common.String(relayAuditLogGroup),
|
|
}}},
|
|
Target: sch.NotificationsTargetDetails{TopicId: &topicID},
|
|
}
|
|
if condition != "" {
|
|
details.Tasks = []sch.TaskDetails{sch.LogRuleTaskDetails{Condition: &condition}}
|
|
}
|
|
if _, err := sc.CreateServiceConnector(ctx, sch.CreateServiceConnectorRequest{
|
|
CreateServiceConnectorDetails: details,
|
|
}); err != nil {
|
|
return RelayResource{}, fmt.Errorf("create service connector: %w", err)
|
|
}
|
|
return waitRelayConnector(ctx, sc, tenancy, name)
|
|
}
|
|
|
|
// reconcileRelayCondition 对齐存量 Connector 的过滤条件:关键事件清单变更
|
|
// (如去除 LaunchInstance)后,已建链路经「一键创建」幂等调用原地更新,无需拆除重建。
|
|
func reconcileRelayCondition(ctx context.Context, sc sch.ServiceConnectorClient, id, condition string) error {
|
|
got, err := sc.GetServiceConnector(ctx, sch.GetServiceConnectorRequest{ServiceConnectorId: &id})
|
|
if err != nil {
|
|
return fmt.Errorf("get service connector: %w", err)
|
|
}
|
|
if !relayConditionDiffers(got.Tasks, condition) {
|
|
return nil
|
|
}
|
|
details := sch.UpdateServiceConnectorDetails{Tasks: []sch.TaskDetails{}}
|
|
if condition != "" {
|
|
details.Tasks = []sch.TaskDetails{sch.LogRuleTaskDetails{Condition: &condition}}
|
|
}
|
|
_, err = sc.UpdateServiceConnector(ctx, sch.UpdateServiceConnectorRequest{
|
|
ServiceConnectorId: &id, UpdateServiceConnectorDetails: details,
|
|
})
|
|
if err != nil {
|
|
return fmt.Errorf("update service connector condition: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// relayConditionDiffers 判断现有任务集与期望条件是否不一致:
|
|
// 期望形态是单条 LogRule 条件;任务数、类型或条件文本不同均视为漂移。
|
|
func relayConditionDiffers(tasks []sch.TaskDetailsResponse, want string) bool {
|
|
if want == "" {
|
|
return len(tasks) > 0
|
|
}
|
|
if len(tasks) != 1 {
|
|
return true
|
|
}
|
|
rule, ok := tasks[0].(sch.LogRuleTaskDetailsResponse)
|
|
if !ok || rule.Condition == nil {
|
|
return true
|
|
}
|
|
return *rule.Condition != want
|
|
}
|
|
|
|
// findRelayConnector 依次按传入的 DisplayName 列表在租户范围内查找存活 Connector;第一个命中即返回。
|
|
func findRelayConnector(ctx context.Context, sc sch.ServiceConnectorClient, tenancy string, names ...string) (RelayResource, bool, error) {
|
|
for _, name := range names {
|
|
list, err := sc.ListServiceConnectors(ctx, sch.ListServiceConnectorsRequest{
|
|
CompartmentId: &tenancy, DisplayName: common.String(name),
|
|
})
|
|
if err != nil {
|
|
return RelayResource{}, false, fmt.Errorf("list service connectors: %w", err)
|
|
}
|
|
for _, item := range list.Items {
|
|
if item.LifecycleState != sch.LifecycleStateDeleted && item.LifecycleState != sch.LifecycleStateDeleting {
|
|
return RelayResource{ID: deref(item.Id), State: string(item.LifecycleState)}, true, nil
|
|
}
|
|
}
|
|
}
|
|
return RelayResource{}, false, nil
|
|
}
|
|
|
|
// waitRelayConnector 轮询新建 Connector 直至 ACTIVE;超时带回已建资源信息。
|
|
// 只按新命名查找(新建资源用的就是新命名)。
|
|
func waitRelayConnector(ctx context.Context, sc sch.ServiceConnectorClient, tenancy, name string) (RelayResource, error) {
|
|
last := RelayResource{Created: true}
|
|
for i := 0; i < relayConnectorPollLimit; i++ {
|
|
select {
|
|
case <-ctx.Done():
|
|
return last, ctx.Err()
|
|
case <-time.After(relayConnectorPollTick):
|
|
}
|
|
res, ok, err := findRelayConnector(ctx, sc, tenancy, name)
|
|
if err != nil {
|
|
return last, err
|
|
}
|
|
if ok {
|
|
last, last.Created = res, true
|
|
if res.State == string(sch.LifecycleStateActive) {
|
|
return last, nil
|
|
}
|
|
}
|
|
}
|
|
return last, fmt.Errorf("service connector not active after %s", time.Duration(relayConnectorPollLimit)*relayConnectorPollTick)
|
|
}
|
|
|
|
// RelayState 实现 Client:聚合链路四资源现状;endpoint 为空时不匹配订阅。
|
|
func (c *RealClient) RelayState(ctx context.Context, cred Credentials, endpoint string) (RelayState, error) {
|
|
var st RelayState
|
|
cp, err := c.onsControlClient(cred)
|
|
if err != nil {
|
|
return st, err
|
|
}
|
|
names := relayResourceNames(cred.TenancyOCID)
|
|
if st.Topic, _, err = findRelayTopic(ctx, cp, cred.TenancyOCID, names.TopicPrefix, legacyRelayTopicPrefix); err != nil {
|
|
return st, err
|
|
}
|
|
if st.Topic.ID != "" {
|
|
dp, err := c.onsDataClient(cred)
|
|
if err != nil {
|
|
return st, err
|
|
}
|
|
if st.Subscription, _, err = findRelaySubscription(ctx, dp, cred.TenancyOCID, st.Topic.ID, endpoint); err != nil {
|
|
return st, err
|
|
}
|
|
}
|
|
return c.relayControlState(ctx, cred, st)
|
|
}
|
|
|
|
// relayControlState 补齐 Connector 与 Policy 两项状态,双路径兼容 legacy 命名。
|
|
func (c *RealClient) relayControlState(ctx context.Context, cred Credentials, st RelayState) (RelayState, error) {
|
|
sc, err := c.schClient(cred)
|
|
if err != nil {
|
|
return st, err
|
|
}
|
|
names := relayResourceNames(cred.TenancyOCID)
|
|
if st.Connector, _, err = findRelayConnector(ctx, sc, cred.TenancyOCID, names.ConnectorName, legacyRelayConnectorName); err != nil {
|
|
return st, err
|
|
}
|
|
ic, err := c.identityClientAt(cred, "")
|
|
if err != nil {
|
|
return st, err
|
|
}
|
|
if st.Policy, _, err = findRelayPolicy(ctx, ic, cred.TenancyOCID, names.PolicyName, legacyRelayPolicyName); err != nil {
|
|
return st, err
|
|
}
|
|
return st, nil
|
|
}
|
|
|
|
// DeleteRelayConnector 实现 Client。
|
|
func (c *RealClient) DeleteRelayConnector(ctx context.Context, cred Credentials, connectorID string) error {
|
|
sc, err := c.schClient(cred)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if _, err := sc.DeleteServiceConnector(ctx, sch.DeleteServiceConnectorRequest{ServiceConnectorId: &connectorID}); err != nil {
|
|
return fmt.Errorf("delete service connector: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// DeleteRelayPolicy 实现 Client;IAM 写操作发往 home region。
|
|
func (c *RealClient) DeleteRelayPolicy(ctx context.Context, cred Credentials, homeRegion, policyID string) error {
|
|
ic, err := c.identityClientAt(cred, homeRegion)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if _, err := ic.DeletePolicy(ctx, identity.DeletePolicyRequest{PolicyId: &policyID}); err != nil {
|
|
return fmt.Errorf("delete policy: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// DeleteRelaySubscription 实现 Client。
|
|
func (c *RealClient) DeleteRelaySubscription(ctx context.Context, cred Credentials, subscriptionID string) error {
|
|
dp, err := c.onsDataClient(cred)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if _, err := dp.DeleteSubscription(ctx, ons.DeleteSubscriptionRequest{SubscriptionId: &subscriptionID}); err != nil {
|
|
return fmt.Errorf("delete ons subscription: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// DeleteRelayTopic 实现 Client。
|
|
func (c *RealClient) DeleteRelayTopic(ctx context.Context, cred Credentials, topicID string) error {
|
|
cp, err := c.onsControlClient(cred)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if _, err := cp.DeleteTopic(ctx, ons.DeleteTopicRequest{TopicId: &topicID}); err != nil {
|
|
return fmt.Errorf("delete ons topic: %w", err)
|
|
}
|
|
return nil
|
|
}
|