Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,10 @@ dropBin:
replicator-service:
air -c configs/air/replicator-service.toml

.PHONY: init-osv-replication
init-osv-replication:
go run ./cmd/first-osv-init-service/main.go

.PHONY: migrate-up
migrate-up:
migrate -path migrations -database $(DB_URL) up
Expand Down
Empty file removed batch.json
Empty file.
3 changes: 3 additions & 0 deletions cmd/api-service/main.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
package main

func main() {}
19 changes: 19 additions & 0 deletions cmd/first-osv-init-service/main.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
package main

import (
"context"
"log"

"github.com/trustpkg/trustpkg-api/db"
"github.com/trustpkg/trustpkg-api/internal/osv"
)

func main() {
db.ConnectDb()

if err := osv.ReplicateEcosystem(context.Background(), osv.FetchOSVEcosystemDumpPayload{
Ecosystem: "npm",
}); err != nil {
log.Printf("OSV replication failed: %v", err)
}
}
9 changes: 5 additions & 4 deletions cmd/replicator-service/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,11 +19,12 @@ func main() {
worker, startWorker := adaptiveWorker.New(ctx)

worker.ScheduleJob(adaptiveWorker.Job{
Handler: func(concurrency int) {
npm.Pipeline(concurrency)
Handler: func(concurrency int, controller *adaptiveWorker.JobController) {
npm.Pipeline(concurrency, controller)

},
SkippedCycles: 0,
ReservationRatio: 1,
SkippedCycles: 29,
ReservationRatio: 0,
})

startWorker()
Expand Down
2 changes: 1 addition & 1 deletion internal/adaptive-worker/constants.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ const (

maxWorkerHeadroom = 0.2

cycleTime = "@every 30s"
cycleTime = "@every 1s"

skipCount = 2

Expand Down
13 changes: 11 additions & 2 deletions internal/adaptive-worker/controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ func (controller *Controller) calculateWorkers() error {
fmt.Println("cpu percent: ", resourceUsage.cpuPercent)
fmt.Println("loadAvg: ", resourceUsage.loadAvg)
fmt.Println("load percent: ", resourceUsage.loadPercent)
fmt.Println("rem percent: ", resourceUsage.ramPercent)
fmt.Println("ram percent: ", resourceUsage.ramPercent)
fmt.Println("resources usage - END-----------------------<")

maxSimultaneousWorkers := getMaxWorkers(resourceUsage.cpuCount)
Expand Down Expand Up @@ -108,7 +108,16 @@ func (controller *Controller) runCycle() {
inferMinSimultaneousWorkers(controller.currentCycle, controller.jobs),
)

currentJob.scheduledJob.Handler(currentJob.concurrency)
if currentJob.canRun {
currentJob.scheduledJob.Handler(currentJob.concurrency, &JobController{
AdjustReservationRatioHandler: func(reservationRatio float64) {
currentJob.scheduledJob.ReservationRatio = reservationRatio
},
AdjustSkippedCyclesHandler: func(skippedCycles int) {
currentJob.scheduledJob.SkippedCycles = skippedCycles
},
})
}
}
}

Expand Down
6 changes: 6 additions & 0 deletions internal/adaptive-worker/job-controler.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
package adaptiveWorker

type JobController struct {
AdjustReservationRatioHandler func(float64)
AdjustSkippedCyclesHandler func(int)
}
2 changes: 1 addition & 1 deletion internal/adaptive-worker/models.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ type calculatedResourcesUsage struct {
ramPercent float64
}

type jobHandler func(concurrency int)
type jobHandler func(concurrency int, controller *JobController)

type Job struct {
Handler jobHandler
Expand Down
3 changes: 3 additions & 0 deletions internal/npm/constants.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,4 +3,7 @@ package npm
const (
npmReplicationUrl = "https://replicate.npmjs.com/registry/_changes"
limit = 10000

npmDownloadUrl = "https://api.npmjs.org/downloads/point"
defaultNpmDownloadPeriod = "last-month"
)
40 changes: 40 additions & 0 deletions internal/npm/http.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package npm

import (
"encoding/json"
"fmt"
"net/http"
"strconv"
)
Expand Down Expand Up @@ -31,3 +32,42 @@ func fetchNpmPackages(since int) (*NpmChangesResponse, error) {

return &packagesData, nil
}

func FetchNpmPackageDownloads(payload FetchNpmPackageDownloadsPayload) (*FetchNpmPackageDownloadsResponse, error) {
var since string
if payload.since == nil {
since = defaultNpmDownloadPeriod
} else {
since = *payload.since
}

url := npmDownloadUrl + "/" + since + "/" + payload.packageName

request, err := http.NewRequest(http.MethodGet, url, nil)
if err != nil {
return nil, err
}

client := &http.Client{}

response, err := client.Do(request)
if err != nil {
return nil, err
}
defer response.Body.Close()

if response.StatusCode < http.StatusOK || response.StatusCode >= http.StatusMultipleChoices {
return nil, fmt.Errorf(
"npm download API returned status %s",
response.Status,
)
}

var packageData FetchNpmPackageDownloadsResponse

if err := json.NewDecoder(response.Body).Decode(&packageData); err != nil {
return nil, err
}

return &packageData, nil
}
12 changes: 12 additions & 0 deletions internal/npm/models.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,3 +23,15 @@ type dbDropPackagesBatchPayload struct {
type dbInsertPackagesBatchPayload struct {
packages []string
}

type FetchNpmPackageDownloadsPayload struct {
packageName string
since *string
}

type FetchNpmPackageDownloadsResponse struct {
Downloads int `json:"downloads"`
End string `json:"end"`
Package string `json:"package"`
Start string `json:"start"`
}
22 changes: 21 additions & 1 deletion internal/npm/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,16 @@ package npm

import (
"context"
"encoding/json"
"fmt"
"os"
"sync"

"github.com/trustpkg/trustpkg-api/db"
adaptiveWorker "github.com/trustpkg/trustpkg-api/internal/adaptive-worker"
)

func Pipeline(concurrency int) error {
func Pipeline(concurrency int, JobController *adaptiveWorker.JobController) error {
ctx := context.Background()

if concurrency < 1 {
Expand Down Expand Up @@ -39,12 +43,28 @@ func Pipeline(concurrency int) error {
if err := savePackagesPage(ctx, changes); err != nil {
errs <- err
}

fmt.Print("---------------------->\n")
enc := json.NewEncoder(os.Stdout)
enc.SetIndent("", " ")
enc.Encode(response.Results)

fmt.Println("response", response.LastSeq)
fmt.Print("----------------------<\n")
}(response.Results)
}

waitGroup.Wait()
close(errs)

data, err := FetchNpmPackageDownloads(FetchNpmPackageDownloadsPayload{
packageName: "next",
})

enc := json.NewEncoder(os.Stdout)
enc.SetIndent("", " ")
enc.Encode(data)

for err := range errs {
if err != nil {
return err
Expand Down
6 changes: 3 additions & 3 deletions internal/npm/sql/insert_packages_batch.sql
Original file line number Diff line number Diff line change
@@ -1,3 +1,3 @@
INSERT INTO packages (name)
SELECT unnest($1::text[])
ON CONFLICT (name) DO NOTHING;
INSERT INTO packages (name, ecosystem)
SELECT unnest($1::text[]), 'npm'
ON CONFLICT (ecosystem, name) DO NOTHING;
9 changes: 9 additions & 0 deletions internal/osv/constants.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
package osv

const (
osvPackageVulnerabilitiesUrl = "https://api.osv.dev/v1/query"
defaultOsvEcosystem = "npm"

npmOsvDumpUrl = "https://storage.googleapis.com/osv-vulnerabilities"
npmOsvDumpUrlSuffix = "all.zip"
)
130 changes: 130 additions & 0 deletions internal/osv/helpers.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,130 @@
package osv

import (
"strconv"
"strings"
"time"
)

type ParsedVulnerability struct {
OSVID string
CVEID *string
Summary *string
Description *string
Severity *string
CVSSScore *float64
CVSSVector *string
PublishedAt *time.Time
ModifiedAt *time.Time

Affected []AffectedPackage
References []string
}

type AffectedPackage struct {
Ecosystem string
PackageName string
IntroducedVersion *string
FixedVersion *string
LastAffectedVersion *string
}

func parseVulnerability(raw OsvVuln) ParsedVulnerability {
parsed := ParsedVulnerability{OSVID: raw.ID, Summary: optionalString(raw.Summary), Description: optionalString(raw.Details), PublishedAt: optionalTime(raw.Published), ModifiedAt: optionalTime(raw.Modified)}
for _, alias := range raw.Aliases {
if strings.HasPrefix(alias, "CVE-") {
parsed.CVEID = optionalString(alias)
break
}
}
for _, severity := range raw.Severity {
if parsed.CVSSVector == nil && severity.Score != "" {
parsed.CVSSVector = optionalString(severity.Score)
}
if score, err := strconv.ParseFloat(severity.Score, 64); err == nil {
parsed.CVSSScore = &score
}
}
if raw.DatabaseSpecific.Severity != "" {
parsed.Severity = optionalString(raw.DatabaseSpecific.Severity)
}
for _, affected := range raw.Affected {
for _, versionRange := range affected.Ranges {
for _, version := range parseAffectedRange(versionRange.Events) {
parsed.Affected = append(parsed.Affected, AffectedPackage{
Ecosystem: affected.Package.Ecosystem, PackageName: affected.Package.Name,
IntroducedVersion: version.introduced,
FixedVersion: version.fixed,
LastAffectedVersion: version.lastAffected,
})
}
}
}
for _, reference := range raw.References {
if reference.URL != "" {
parsed.References = append(parsed.References, reference.URL)
}
}
return parsed
}

type affectedVersion struct {
introduced *string
fixed *string
lastAffected *string
}

func parseAffectedRange(events []OsvEvent) []affectedVersion {
versions := make([]affectedVersion, 0, len(events))
var introduced *string

for _, event := range events {
if event.Introduced != "" {
if introduced != nil {
versions = append(versions, affectedVersion{introduced: introduced})
}
introduced = optionalString(event.Introduced)
}

if event.Fixed != "" {
if introduced == nil {
introduced = optionalString("0")
}
versions = append(versions, affectedVersion{
introduced: introduced,
fixed: optionalString(event.Fixed),
})
introduced = nil
}

if event.LastAffected != "" {
if introduced == nil {
introduced = optionalString("0")
}
versions = append(versions, affectedVersion{
introduced: introduced,
lastAffected: optionalString(event.LastAffected),
})
introduced = nil
}
}

if introduced != nil {
versions = append(versions, affectedVersion{introduced: introduced})
}
return versions
}

func optionalString(value string) *string {
if value == "" {
return nil
}
return &value
}

func optionalTime(value time.Time) *time.Time {
if value.IsZero() {
return nil
}
return &value
}
Loading