Skip to content

Commit

Permalink
refactor: refact the notification job API and life process
Browse files Browse the repository at this point in the history
1. Introduce new APIs for webhook jobs management.
2. Refact legacy APIs for backforward compatible.
3. Migrate the webhook jobs process to unified execution/task framework.

Closes: goharbor#18210

Signed-off-by: chlins <chenyuzh@vmware.com>
  • Loading branch information
chlins committed Feb 21, 2023
1 parent 99b3711 commit 45f87c7
Show file tree
Hide file tree
Showing 27 changed files with 1,754 additions and 1,087 deletions.
150 changes: 145 additions & 5 deletions api/v2.0/swagger.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -2665,8 +2665,128 @@ paths:
$ref: '#/responses/404'
'500':
$ref: '#/responses/500'
'/projects/{project_name_or_id}/webhook/policies/{webhook_policy_id}/executions':
get:
summary: List executions for a specific webhook policy
description: |
This endpoint returns the executions of a specific webhook policy.
tags:
- webhook
operationId: ListExecutionsOfWebhookPolicy
parameters:
- $ref: '#/parameters/requestId'
- $ref: '#/parameters/isResourceName'
- $ref: '#/parameters/projectNameOrId'
- $ref: '#/parameters/webhookPolicyId'
- $ref: '#/parameters/page'
- $ref: '#/parameters/pageSize'
- $ref: '#/parameters/query'
- $ref: '#/parameters/sort'
responses:
'200':
description: List webhook executions success
headers:
X-Total-Count:
description: The total count of executions
type: integer
Link:
description: Link refers to the previous page and next page
type: string
schema:
type: array
items:
$ref: '#/definitions/Execution'
'400':
$ref: '#/responses/400'
'401':
$ref: '#/responses/401'
'404':
$ref: '#/responses/404'
'403':
$ref: '#/responses/403'
'500':
$ref: '#/responses/500'
'/projects/{project_name_or_id}/webhook/policies/{webhook_policy_id}/executions/{execution_id}/tasks':
get:
summary: List tasks for a specific webhook execution
description: |
This endpoint returns the tasks of a specific webhook execution.
tags:
- webhook
operationId: ListTasksOfWebhookExecution
parameters:
- $ref: '#/parameters/requestId'
- $ref: '#/parameters/isResourceName'
- $ref: '#/parameters/projectNameOrId'
- $ref: '#/parameters/webhookPolicyId'
- $ref: '#/parameters/executionId'
- $ref: '#/parameters/page'
- $ref: '#/parameters/pageSize'
- $ref: '#/parameters/query'
- $ref: '#/parameters/sort'
responses:
'200':
description: List tasks of webhook executions success
headers:
X-Total-Count:
description: The total count of tasks
type: integer
Link:
description: Link refers to the previous page and next page
type: string
schema:
type: array
items:
$ref: '#/definitions/Task'
'400':
$ref: '#/responses/400'
'401':
$ref: '#/responses/401'
'404':
$ref: '#/responses/404'
'403':
$ref: '#/responses/403'
'500':
$ref: '#/responses/500'
'/projects/{project_name_or_id}/webhook/policies/{webhook_policy_id}/executions/{execution_id}/tasks/{task_id}/log':
get:
summary: Get logs for a specific webhook task
description: |
This endpoint returns the logs of a specific webhook task.
tags:
- webhook
operationId: GetLogsOfWebhookTask
produces:
- text/plain
parameters:
- $ref: '#/parameters/requestId'
- $ref: '#/parameters/isResourceName'
- $ref: '#/parameters/projectNameOrId'
- $ref: '#/parameters/webhookPolicyId'
- $ref: '#/parameters/executionId'
- $ref: '#/parameters/taskId'
responses:
'200':
description: Get log success
headers:
Content-Type:
description: Content type of response
type: string
schema:
type: string
'400':
$ref: '#/responses/400'
'401':
$ref: '#/responses/401'
'404':
$ref: '#/responses/404'
'403':
$ref: '#/responses/403'
'500':
$ref: '#/responses/500'
'/projects/{project_name_or_id}/webhook/lasttrigger':
get:
deprecated: true
summary: Get project webhook policy last trigger info
description: |
This endpoint returns last trigger information of project webhook policy.
Expand Down Expand Up @@ -2694,6 +2814,7 @@ paths:
$ref: '#/responses/500'
'/projects/{project_name_or_id}/webhook/jobs':
get:
deprecated: true
summary: List project webhook jobs
description: |
This endpoint returns webhook jobs of a project.
Expand Down Expand Up @@ -2746,7 +2867,7 @@ paths:
'/projects/{project_name_or_id}/webhook/events':
get:
summary: Get supported event types and notify types.
description: Get supportted event types and notify types.
description: Get supported event types and notify types.
tags:
- webhook
operationId: GetSupportedEventTypes
Expand Down Expand Up @@ -8445,7 +8566,7 @@ definitions:
description: 'The group type, 1 for LDAP group, 2 for HTTP group, 3 for OIDC group.'
SupportedWebhookEventTypes:
type: object
description: Supportted webhook event types and notify types.
description: Supported webhook event types and notify types.
properties:
event_type:
type: array
Expand All @@ -8455,14 +8576,33 @@ definitions:
type: array
items:
$ref: '#/definitions/NotifyType'
payload_formats:
type: array
items:
$ref: '#/definitions/PayloadFormat'
EventType:
type: string
description: Webhook supportted event type.
example: 'pullImage'
description: Webhook supported event type.
example: 'PULL_ARTIFACT'
NotifyType:
type: string
description: Webhook supportted notify type.
description: Webhook supported notify type.
example: 'http'
PayloadFormatType:
type: string
description: The type of webhook paylod format.
example: 'cloudevent'
PayloadFormat:
type: object
description: Webhook supported payload format type collections.
properties:
notify_type:
$ref: '#/definitions/NotifyType'
formats:
type: array
description: The supported payload formats for this notify type.
items:
$ref: '#/definitions/PayloadFormatType'

WebhookTargetObject:
type: object
Expand Down
195 changes: 195 additions & 0 deletions src/controller/webhook/controller.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,195 @@
// Copyright Project Harbor Authors
//
// Licensed 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 webhook

import (
"context"
"time"

"github.com/goharbor/harbor/src/jobservice/job"
"github.com/goharbor/harbor/src/lib/errors"
"github.com/goharbor/harbor/src/lib/q"
"github.com/goharbor/harbor/src/pkg/notification/policy"
"github.com/goharbor/harbor/src/pkg/notification/policy/model"
"github.com/goharbor/harbor/src/pkg/task"
)

var (
// Ctl is a global webhook controller instance
Ctl = NewController()

// webhookJobVendors represents webhook(http) or slack.
webhookJobVendors = q.NewOrList([]interface{}{job.WebhookJobVendorType, job.SlackJobVendorType})
)

type Controller interface {
// CreatePolicy creates webhook policy
CreatePolicy(ctx context.Context, policy *model.Policy) (int64, error)
// ListPolicies lists webhook policies filter by query
ListPolicies(ctx context.Context, query *q.Query) ([]*model.Policy, error)
// CountPolicies counts webhook policies filter by query
CountPolicies(ctx context.Context, query *q.Query) (int64, error)
// GetPolicy gets webhook policy by specified ID
GetPolicy(ctx context.Context, id int64) (*model.Policy, error)
// UpdatePolicy updates webhook policy
UpdatePolicy(ctx context.Context, policy *model.Policy) error
// DeletePolicy deletes webhook policy by specified ID
DeletePolicy(ctx context.Context, policyID int64) error
// GetRelatedPolices gets related policies by the input project id and event type
GetRelatedPolices(ctx context.Context, projectID int64, eventType string) ([]*model.Policy, error)

// CountExecutions counts executions under the webhook policy
CountExecutions(ctx context.Context, policyID int64, query *q.Query) (int64, error)
// ListExecutions lists executions under the webhook policy
ListExecutions(ctx context.Context, policyID int64, query *q.Query) ([]*task.Execution, error)
// CountTasks counts tasks under the webhook execution
CountTasks(ctx context.Context, execID int64, query *q.Query) (int64, error)
// ListTasks lists tasks under the webhook execution
ListTasks(ctx context.Context, execID int64, query *q.Query) ([]*task.Task, error)
// GetTask gets the webhook task by the specified ID
GetTask(ctx context.Context, taskID int64) (*task.Task, error)
// GetTaskLog gets task log
GetTaskLog(ctx context.Context, taskID int64) ([]byte, error)

// GetLastTriggerTime gets policy last trigger time group by event type
GetLastTriggerTime(ctx context.Context, eventType string, policyID int64) (time.Time, error)
}

type controller struct {
policyMgr policy.Manager
execMgr task.ExecutionManager
taskMgr task.Manager
}

func NewController() Controller {
return &controller{
policyMgr: policy.Mgr,
execMgr: task.ExecMgr,
taskMgr: task.Mgr,
}
}

func (c *controller) CreatePolicy(ctx context.Context, policy *model.Policy) (int64, error) {
return c.policyMgr.Create(ctx, policy)
}

func (c *controller) ListPolicies(ctx context.Context, query *q.Query) ([]*model.Policy, error) {
return c.policyMgr.List(ctx, query)
}

func (c *controller) CountPolicies(ctx context.Context, query *q.Query) (int64, error) {
return c.policyMgr.Count(ctx, query)
}

func (c *controller) GetPolicy(ctx context.Context, id int64) (*model.Policy, error) {
return c.policyMgr.Get(ctx, id)
}

func (c *controller) UpdatePolicy(ctx context.Context, policy *model.Policy) error {
return c.policyMgr.Update(ctx, policy)
}

func (c *controller) DeletePolicy(ctx context.Context, policyID int64) error {
// delete executions under the webhook policy,
// there are two vendor types(webhook & slack) needs to be deleted.
if err := c.execMgr.DeleteByVendor(ctx, job.WebhookJobVendorType, policyID); err != nil {
return errors.Wrapf(err, "failed to delete executions for webhook of policy %d", policyID)
}
if err := c.execMgr.DeleteByVendor(ctx, job.SlackJobVendorType, policyID); err != nil {
return errors.Wrapf(err, "failed to delete executions for slack of policy %d", policyID)
}

return c.policyMgr.Delete(ctx, policyID)
}

func (c *controller) GetRelatedPolices(ctx context.Context, projectID int64, eventType string) ([]*model.Policy, error) {
return c.policyMgr.GetRelatedPolices(ctx, projectID, eventType)
}

func (c *controller) CountExecutions(ctx context.Context, policyID int64, query *q.Query) (int64, error) {
return c.execMgr.Count(ctx, buildExecutionQuery(policyID, query))
}

func (c *controller) ListExecutions(ctx context.Context, policyID int64, query *q.Query) ([]*task.Execution, error) {
return c.execMgr.List(ctx, buildExecutionQuery(policyID, query))
}

func (c *controller) CountTasks(ctx context.Context, execID int64, query *q.Query) (int64, error) {
return c.taskMgr.Count(ctx, buildTaskQuery(execID, query))
}

func (c *controller) ListTasks(ctx context.Context, execID int64, query *q.Query) ([]*task.Task, error) {
return c.taskMgr.List(ctx, buildTaskQuery(execID, query))
}

func (c *controller) GetTask(ctx context.Context, taskID int64) (*task.Task, error) {
query := q.New(q.KeyWords{
"id": taskID,
"vendor_type": webhookJobVendors,
})
tasks, err := c.taskMgr.List(ctx, query)
if err != nil {
return nil, err
}

if len(tasks) == 0 {
return nil, errors.New(nil).WithCode(errors.NotFoundCode).
WithMessage("webhook task %d not found", taskID)
}
return tasks[0], nil
}

func (c *controller) GetTaskLog(ctx context.Context, taskID int64) ([]byte, error) {
// ensure the webhook task exist
_, err := c.GetTask(ctx, taskID)
if err != nil {
return nil, err
}

return c.taskMgr.GetLog(ctx, taskID)
}

func buildExecutionQuery(policyID int64, query *q.Query) *q.Query {
query = q.MustClone(query)
query.Keywords["vendor_type"] = webhookJobVendors
query.Keywords["vendor_id"] = policyID
return query
}

func buildTaskQuery(execID int64, query *q.Query) *q.Query {
query = q.MustClone(query)
query.Keywords["vendor_type"] = webhookJobVendors
query.Keywords["execution_id"] = execID
return query
}

func (c *controller) GetLastTriggerTime(ctx context.Context, eventType string, policyID int64) (time.Time, error) {
query := q.New(q.KeyWords{
"vendor_type": webhookJobVendors,
"vendor_id": policyID,
"ExtraAttrs.type": eventType,
})
// fetch the latest execution sort by start_time
execs, err := c.execMgr.List(ctx, query.First(q.NewSort("start_time", true)))
if err != nil {
return time.Time{}, err
}

if len(execs) > 0 {
return execs[0].StartTime, nil
}

return time.Time{}, nil
}
Loading

0 comments on commit 45f87c7

Please sign in to comment.