ci / go (push) Waiting to run
ci / go-db (agent) (push) Waiting to run
ci / go-db (config) (push) Waiting to run
ci / go-db (db) (push) Waiting to run
ci / go-db (evidence) (push) Waiting to run
ci / go-db (llmrec) (push) Waiting to run
ci / go-db (server) (push) Waiting to run
detections / detections (push) Waiting to run
web / web (push) Waiting to run
docs / links (push) Canceled after 0s
536 lines
16 KiB
Go
536 lines
16 KiB
Go
package server
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"net/http"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/Autumn-27/artex/db"
|
|
"github.com/Autumn-27/artex/mcphttp"
|
|
)
|
|
|
|
// 자산 동기화(ScopeSentry 데이터 소스).
|
|
//
|
|
// ScopeSentry 는 ASM 자산 측량 플랫폼으로, MCP 인터페이스를 통해 [프로젝트]/[작업] 두 차원으로
|
|
// 서브도메인·웹 애플리케이션·서비스 등의 자산을 가져와 ARTEX 의 회사 + 자산 모델로 매핑한다.
|
|
// 데이터 소스 자체는 "ScopeSentry" 라는 이름의 http 전송 MCP 행이다(url + X-API-Key 헤더를 mcp_servers 에 저장).
|
|
//
|
|
// 에이전트 도구 계층과 달리, 여기서는 mcphttp.Client.Call 로 MCP 도구를 직접 호출해 원본 JSON 을 받고,
|
|
// AskUser 권한 팝업을 띄우지 않는다(백그라운드 일괄 동기화).
|
|
|
|
const (
|
|
scopeSentryMCPName = "ScopeSentry"
|
|
syncMaxPerType = 5000 // 단일 대상·단일 유형당 저장 보호 상한
|
|
syncDefaultPage = 100
|
|
)
|
|
|
|
// 사용자 노출 메시지(writeErr 응답 + 동기화 결과 warnings/errors JSON). 식별자·enum·MCP 도구명·
|
|
// JSON 필드명(dimension·targets 등)은 원문 그대로 두고, 중국어 문구만 한국어로 옮긴다(BRIEF F3 방침).
|
|
// UI(web/messages/ko.json 의 sync 네임스페이스)와 "데이터 소스" 표기를 맞춘다.
|
|
const (
|
|
errSSDataSourceMissingFmt = "데이터 소스 %s 가 없습니다. 먼저 생성해 주세요"
|
|
errSSDataSourceNoURLFmt = "데이터 소스 %s 에 URL 이 설정되지 않았습니다. 먼저 설정해 주세요"
|
|
errSSListProjectsPrefix = "list_projects_data 호출에 실패했습니다: "
|
|
errSSParseProjectsPrefix = "프로젝트 목록을 해석하지 못했습니다: "
|
|
errSSListTasksPrefix = "list_tasks 호출에 실패했습니다: "
|
|
errSSParseTasksPrefix = "작업 목록을 해석하지 못했습니다: "
|
|
errSSDimension = "dimension 값은 project 또는 task 중 하나여야 합니다"
|
|
errSSTargetsEmpty = "targets 항목은 비워 둘 수 없습니다"
|
|
warnSSProjectMetaFmt = "프로젝트 %s 의 상세 정보를 가져오지 못했습니다: %v"
|
|
warnSSCompanyCreateFmt = "회사 %s 생성에 실패했습니다: %v"
|
|
warnSSUnknownAssetPrefix = "알 수 없는 자산 유형이라 건너뛰었습니다: "
|
|
errSSFetchFmt = "%s(%s) 자산을 가져오지 못했습니다: %v"
|
|
warnSSTruncatedFmt = "%s(%s) 가 %d 건 상한에 도달해 잘렸습니다"
|
|
errSSParseSubdomainPrefix = "subdomain 응답을 해석하지 못했습니다: "
|
|
errSSParseAppPrefix = "app 응답을 해석하지 못했습니다: "
|
|
errSSParseServicePrefix = "service 응답을 해석하지 못했습니다: "
|
|
)
|
|
|
|
// findMCPByName returns the MCP server row with the given name, or nil.
|
|
func (s *Server) findMCPByName(name string) (*db.MCPServer, error) {
|
|
all, err := s.m.pg.ListMCP()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
for _, m := range all {
|
|
if m.Name == name {
|
|
return m, nil
|
|
}
|
|
}
|
|
return nil, nil
|
|
}
|
|
|
|
// scopeSentryClient dials the configured ScopeSentry MCP (http transport). The env
|
|
// map doubles as HTTP headers (X-API-Key). Callers must Close the client.
|
|
func (s *Server) scopeSentryClient(ctx context.Context) (*mcphttp.Client, error) {
|
|
m, err := s.findMCPByName(scopeSentryMCPName)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if m == nil {
|
|
return nil, fmt.Errorf(errSSDataSourceMissingFmt, scopeSentryMCPName)
|
|
}
|
|
if m.URL == "" {
|
|
return nil, fmt.Errorf(errSSDataSourceNoURLFmt, scopeSentryMCPName)
|
|
}
|
|
return mcphttp.New(ctx, m.Name, m.URL, jsonStrMap(m.Env), m.Insecure)
|
|
}
|
|
|
|
// envHasValue reports whether the env/header map has any non-empty value
|
|
// (i.e. an API key / Authorization header was filled in).
|
|
func envHasValue(raw json.RawMessage) bool {
|
|
for _, v := range jsonStrMap(raw) {
|
|
if strings.TrimSpace(v) != "" {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// ---------- GET /api/sync/scopesentry/status ----------
|
|
|
|
func (s *Server) syncSSStatus(w http.ResponseWriter, r *http.Request) {
|
|
if s.pg(w) == nil {
|
|
return
|
|
}
|
|
m, err := s.findMCPByName(scopeSentryMCPName)
|
|
if err != nil {
|
|
writeErr(w, 500, err.Error())
|
|
return
|
|
}
|
|
resp := map[string]any{
|
|
"exists": m != nil,
|
|
"configured": false,
|
|
"enabled": false,
|
|
"reachable": false,
|
|
"tools": []string{},
|
|
}
|
|
if m == nil {
|
|
writeJSON(w, 200, resp)
|
|
return
|
|
}
|
|
configured := m.URL != "" && envHasValue(m.Env)
|
|
resp["configured"] = configured
|
|
resp["enabled"] = m.Enabled
|
|
resp["url"] = m.URL
|
|
if m.Tools != nil {
|
|
resp["tools"] = m.Tools
|
|
}
|
|
// Light reachability probe only when it can actually connect.
|
|
if configured && m.Enabled {
|
|
ctx, cancel := context.WithTimeout(r.Context(), 12*time.Second)
|
|
defer cancel()
|
|
if cl, cerr := mcphttp.New(ctx, m.Name, m.URL, jsonStrMap(m.Env), m.Insecure); cerr == nil {
|
|
if _, terr := cl.Tools(ctx); terr == nil {
|
|
resp["reachable"] = true
|
|
}
|
|
_ = cl.Close()
|
|
}
|
|
}
|
|
writeJSON(w, 200, resp)
|
|
}
|
|
|
|
// ---------- POST /api/sync/scopesentry/datasource ----------
|
|
// Creates the placeholder row if missing, and/or sets URL + API key and enables it.
|
|
|
|
func (s *Server) syncSSDatasource(w http.ResponseWriter, r *http.Request) {
|
|
pg := s.pg(w)
|
|
if pg == nil {
|
|
return
|
|
}
|
|
var body struct {
|
|
URL string `json:"url"`
|
|
APIKey string `json:"api_key"`
|
|
}
|
|
_ = decode(r, &body) // empty body = just create the placeholder
|
|
body.URL = strings.TrimSpace(body.URL)
|
|
body.APIKey = strings.TrimSpace(body.APIKey)
|
|
|
|
m, err := s.findMCPByName(scopeSentryMCPName)
|
|
if err != nil {
|
|
writeErr(w, 500, err.Error())
|
|
return
|
|
}
|
|
if m == nil {
|
|
m = &db.MCPServer{Name: scopeSentryMCPName, Transport: "http", Env: json.RawMessage(`{"X-API-Key":""}`)}
|
|
}
|
|
if body.URL != "" {
|
|
m.URL = body.URL
|
|
}
|
|
if body.APIKey != "" {
|
|
env, _ := json.Marshal(map[string]string{"X-API-Key": body.APIKey})
|
|
m.Env = env
|
|
}
|
|
// Enable only once it can actually be used.
|
|
m.Enabled = m.URL != "" && envHasValue(m.Env)
|
|
|
|
id, err := pg.SaveMCP(m)
|
|
if err != nil {
|
|
writeErr(w, 500, err.Error())
|
|
return
|
|
}
|
|
m.ID = id
|
|
// Best-effort tool discovery so the status card shows tools right away.
|
|
if m.Enabled {
|
|
ctx, cancel := context.WithTimeout(r.Context(), 30*time.Second)
|
|
_ = s.discoverAndCacheMCP(ctx, m)
|
|
cancel()
|
|
}
|
|
writeJSON(w, 200, map[string]any{"id": id, "enabled": m.Enabled})
|
|
}
|
|
|
|
// ---------- GET /api/sync/scopesentry/projects ----------
|
|
|
|
func (s *Server) syncSSProjects(w http.ResponseWriter, r *http.Request) {
|
|
if s.pg(w) == nil {
|
|
return
|
|
}
|
|
ctx, cancel := context.WithTimeout(r.Context(), 60*time.Second)
|
|
defer cancel()
|
|
cl, err := s.scopeSentryClient(ctx)
|
|
if err != nil {
|
|
writeErr(w, 502, err.Error())
|
|
return
|
|
}
|
|
defer cl.Close()
|
|
|
|
args := map[string]any{
|
|
"pageIndex": queryInt(r, "page", 1),
|
|
"pageSize": queryInt(r, "size", 50),
|
|
}
|
|
if q := r.URL.Query().Get("search"); q != "" {
|
|
args["search"] = q
|
|
}
|
|
text, err := cl.Call(ctx, "list_projects_data", args)
|
|
if err != nil {
|
|
writeErr(w, 502, errSSListProjectsPrefix+err.Error())
|
|
return
|
|
}
|
|
// {result:{All:[{id,name,logo,AssetCount,tag}], <tag>:[...]}, tag:{...}}
|
|
var env struct {
|
|
Result map[string]json.RawMessage `json:"result"`
|
|
Tag map[string]int `json:"tag"`
|
|
}
|
|
if err := json.Unmarshal([]byte(text), &env); err != nil {
|
|
writeErr(w, 502, errSSParseProjectsPrefix+err.Error())
|
|
return
|
|
}
|
|
projects := json.RawMessage("[]")
|
|
if raw, ok := env.Result["All"]; ok && len(raw) > 0 {
|
|
projects = raw
|
|
}
|
|
writeJSON(w, 200, map[string]any{"projects": projects, "tag": env.Tag})
|
|
}
|
|
|
|
// ---------- GET /api/sync/scopesentry/tasks ----------
|
|
|
|
func (s *Server) syncSSTasks(w http.ResponseWriter, r *http.Request) {
|
|
if s.pg(w) == nil {
|
|
return
|
|
}
|
|
ctx, cancel := context.WithTimeout(r.Context(), 60*time.Second)
|
|
defer cancel()
|
|
cl, err := s.scopeSentryClient(ctx)
|
|
if err != nil {
|
|
writeErr(w, 502, err.Error())
|
|
return
|
|
}
|
|
defer cl.Close()
|
|
|
|
args := map[string]any{
|
|
"pageIndex": queryInt(r, "page", 1),
|
|
"pageSize": queryInt(r, "size", 50),
|
|
}
|
|
if q := r.URL.Query().Get("search"); q != "" {
|
|
args["search"] = q
|
|
}
|
|
text, err := cl.Call(ctx, "list_tasks", args)
|
|
if err != nil {
|
|
writeErr(w, 502, errSSListTasksPrefix+err.Error())
|
|
return
|
|
}
|
|
var env struct {
|
|
List json.RawMessage `json:"list"`
|
|
}
|
|
if err := json.Unmarshal([]byte(text), &env); err != nil {
|
|
writeErr(w, 502, errSSParseTasksPrefix+err.Error())
|
|
return
|
|
}
|
|
tasks := env.List
|
|
if len(tasks) == 0 {
|
|
tasks = json.RawMessage("[]")
|
|
}
|
|
writeJSON(w, 200, map[string]any{"tasks": tasks})
|
|
}
|
|
|
|
// ---------- POST /api/sync/scopesentry/sync ----------
|
|
|
|
type ssSyncReq struct {
|
|
Dimension string `json:"dimension"` // "project" | "task"
|
|
Targets []string `json:"targets"` // project ObjectIDs, or task names
|
|
AssetTypes []string `json:"asset_types"` // subdomain | app | service
|
|
CreateCompany bool `json:"create_company"`
|
|
PageSize int `json:"page_size"`
|
|
}
|
|
|
|
func (s *Server) syncSSRun(w http.ResponseWriter, r *http.Request) {
|
|
if s.pg(w) == nil {
|
|
return
|
|
}
|
|
as := s.assetStore()
|
|
cs := s.companyStore()
|
|
if as == nil || cs == nil {
|
|
writeErr(w, 503, "database unavailable")
|
|
return
|
|
}
|
|
var req ssSyncReq
|
|
if err := decode(r, &req); err != nil {
|
|
writeErr(w, 400, "invalid JSON: "+err.Error())
|
|
return
|
|
}
|
|
if req.Dimension != "project" && req.Dimension != "task" {
|
|
writeErr(w, 400, errSSDimension)
|
|
return
|
|
}
|
|
if len(req.Targets) == 0 {
|
|
writeErr(w, 400, errSSTargetsEmpty)
|
|
return
|
|
}
|
|
if len(req.AssetTypes) == 0 {
|
|
req.AssetTypes = []string{"subdomain", "service", "app"}
|
|
}
|
|
pageSize := req.PageSize
|
|
if pageSize <= 0 || pageSize > 500 {
|
|
pageSize = syncDefaultPage
|
|
}
|
|
|
|
ctx, cancel := context.WithTimeout(r.Context(), 10*time.Minute)
|
|
defer cancel()
|
|
cl, err := s.scopeSentryClient(ctx)
|
|
if err != nil {
|
|
writeErr(w, 502, err.Error())
|
|
return
|
|
}
|
|
defer cl.Close()
|
|
|
|
synced := map[string]int{"subdomain": 0, "app": 0, "service": 0, "ip": 0}
|
|
var companies []string
|
|
var warnings, errs []string
|
|
madeCompany := false
|
|
|
|
for _, target := range req.Targets {
|
|
filter := map[string]any{}
|
|
switch req.Dimension {
|
|
case "project":
|
|
filter["project"] = []string{target}
|
|
if req.CreateCompany {
|
|
name, roots, perr := s.ssProjectMeta(ctx, cl, target)
|
|
if perr != nil {
|
|
warnings = append(warnings, fmt.Sprintf(warnSSProjectMetaFmt, target, perr))
|
|
} else if name != "" {
|
|
cid, cerr := cs.UpsertByName(name)
|
|
if cerr != nil {
|
|
warnings = append(warnings, fmt.Sprintf(warnSSCompanyCreateFmt, name, cerr))
|
|
} else {
|
|
companies = append(companies, name)
|
|
madeCompany = true
|
|
if len(roots) > 0 {
|
|
cs.AddScope(cid, roots, "scopesentry:"+name)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
case "task":
|
|
filter["task"] = []string{target}
|
|
}
|
|
|
|
for _, at := range req.AssetTypes {
|
|
ssType, ok := map[string]string{"subdomain": "subdomain", "service": "asset", "app": "app"}[at]
|
|
if !ok {
|
|
warnings = append(warnings, warnSSUnknownAssetPrefix+at)
|
|
continue
|
|
}
|
|
items, truncated, ferr := s.ssPageAll(ctx, cl, ssType, filter, pageSize)
|
|
if ferr != nil {
|
|
errs = append(errs, fmt.Sprintf(errSSFetchFmt, at, target, ferr))
|
|
continue
|
|
}
|
|
if truncated {
|
|
warnings = append(warnings, fmt.Sprintf(warnSSTruncatedFmt, at, target, syncMaxPerType))
|
|
}
|
|
for _, raw := range items {
|
|
if e := s.ssIngest(as, at, raw, synced); e != "" {
|
|
errs = append(errs, e)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// Rebuild company attribution so freshly-synced assets attach to their company
|
|
// (scope was written above; assets came after).
|
|
if madeCompany {
|
|
_ = cs.RecomputeAttribution()
|
|
}
|
|
|
|
writeJSON(w, 200, map[string]any{
|
|
"synced": synced,
|
|
"companies": companies,
|
|
"warnings": warnings,
|
|
"errors": errs,
|
|
})
|
|
}
|
|
|
|
// ssProjectMeta fetches a project's display name and its target root domains
|
|
// (get_project.target is a newline-separated list) for company + scope creation.
|
|
func (s *Server) ssProjectMeta(ctx context.Context, cl *mcphttp.Client, projectID string) (name string, roots []string, err error) {
|
|
text, err := cl.Call(ctx, "get_project", map[string]any{"id": projectID})
|
|
if err != nil {
|
|
return "", nil, err
|
|
}
|
|
var p struct {
|
|
Name string `json:"name"`
|
|
Target string `json:"target"`
|
|
}
|
|
if err := json.Unmarshal([]byte(text), &p); err != nil {
|
|
return "", nil, err
|
|
}
|
|
for _, line := range strings.Split(p.Target, "\n") {
|
|
if t := strings.TrimSpace(line); t != "" {
|
|
roots = append(roots, t)
|
|
}
|
|
}
|
|
return p.Name, roots, nil
|
|
}
|
|
|
|
// ssPageAll pages through list_assets for one asset_type + filter until a short/empty
|
|
// page or the protection cap. The {list:[...]} envelope has no total, so we stop when
|
|
// a page returns fewer than pageSize items.
|
|
func (s *Server) ssPageAll(ctx context.Context, cl *mcphttp.Client, ssType string, filter map[string]any, pageSize int) (items []json.RawMessage, truncated bool, err error) {
|
|
for page := 1; ; page++ {
|
|
args := map[string]any{"asset_type": ssType, "pageIndex": page, "pageSize": pageSize}
|
|
if len(filter) > 0 {
|
|
args["filter"] = filter
|
|
}
|
|
text, cerr := cl.Call(ctx, "list_assets", args)
|
|
if cerr != nil {
|
|
return items, truncated, cerr
|
|
}
|
|
var env struct {
|
|
List []json.RawMessage `json:"list"`
|
|
}
|
|
if uerr := json.Unmarshal([]byte(text), &env); uerr != nil {
|
|
return items, truncated, uerr
|
|
}
|
|
if len(env.List) == 0 {
|
|
break
|
|
}
|
|
items = append(items, env.List...)
|
|
if len(items) >= syncMaxPerType {
|
|
items = items[:syncMaxPerType]
|
|
truncated = true
|
|
break
|
|
}
|
|
if len(env.List) < pageSize {
|
|
break
|
|
}
|
|
}
|
|
return items, truncated, nil
|
|
}
|
|
|
|
// ssIngest maps one ScopeSentry asset JSON to the ARTEX asset store and upserts it.
|
|
// Returns a non-empty error string on failure. synced is incremented per kind.
|
|
func (s *Server) ssIngest(as *db.AssetStore, assetType string, raw json.RawMessage, synced map[string]int) string {
|
|
switch assetType {
|
|
case "subdomain":
|
|
var it struct {
|
|
Host string `json:"host"`
|
|
Type string `json:"type"`
|
|
Value []string `json:"value"`
|
|
IP []string `json:"ip"`
|
|
}
|
|
if err := json.Unmarshal(raw, &it); err != nil {
|
|
return errSSParseSubdomainPrefix + err.Error()
|
|
}
|
|
if it.Host == "" {
|
|
return ""
|
|
}
|
|
if _, err := as.UpsertSubdomain(db.UpsertSubdomainReq{Domain: it.Host, RecordType: it.Type, RecordValue: it.Value}); err != nil {
|
|
return "subdomain " + it.Host + ": " + err.Error()
|
|
}
|
|
synced["subdomain"]++
|
|
for _, ip := range it.IP {
|
|
if ip != "" {
|
|
if _, err := as.UpsertIP(db.UpsertIPReq{IP: ip, BoundDomains: []string{it.Host}}); err == nil {
|
|
synced["ip"]++
|
|
}
|
|
}
|
|
}
|
|
case "app":
|
|
var it struct {
|
|
Name string `json:"name"`
|
|
Category string `json:"category"`
|
|
Description string `json:"description"`
|
|
ICP string `json:"icp"`
|
|
}
|
|
if err := json.Unmarshal(raw, &it); err != nil {
|
|
return errSSParseAppPrefix + err.Error()
|
|
}
|
|
if it.Name == "" {
|
|
return ""
|
|
}
|
|
if _, err := as.UpsertApp(db.UpsertAppReq{Name: it.Name, Category: it.Category, Description: it.Description, ICP: it.ICP}); err != nil {
|
|
return "app " + it.Name + ": " + err.Error()
|
|
}
|
|
synced["app"]++
|
|
case "service":
|
|
var it struct {
|
|
Domain string `json:"domain"`
|
|
IP string `json:"ip"`
|
|
Port string `json:"port"`
|
|
Service string `json:"service"`
|
|
URL string `json:"url"`
|
|
Title string `json:"title"`
|
|
Status *int `json:"status"`
|
|
Products []string `json:"products"`
|
|
Icon string `json:"icon"`
|
|
}
|
|
if err := json.Unmarshal(raw, &it); err != nil {
|
|
return errSSParseServicePrefix + err.Error()
|
|
}
|
|
if it.Service == "http" || it.URL != "" {
|
|
if it.URL == "" {
|
|
return ""
|
|
}
|
|
if _, err := as.UpsertHTTPService(db.UpsertHTTPServiceReq{
|
|
URL: it.URL, Technologies: it.Products, StatusCode: it.Status,
|
|
PageTitle: it.Title, FaviconMMH3: it.Icon, IP: it.IP,
|
|
}); err != nil {
|
|
return "service " + it.URL + ": " + err.Error()
|
|
}
|
|
} else {
|
|
port, _ := strconv.Atoi(it.Port)
|
|
if _, err := as.UpsertOtherService(db.UpsertOtherServiceReq{
|
|
Domain: it.Domain, IP: it.IP, Port: port, ServiceName: it.Service,
|
|
}); err != nil {
|
|
return "service " + it.IP + ":" + it.Port + ": " + err.Error()
|
|
}
|
|
}
|
|
synced["service"]++
|
|
}
|
|
return ""
|
|
}
|
|
|
|
// queryInt reads an int query param with a default fallback.
|
|
func queryInt(r *http.Request, key string, def int) int {
|
|
if v := r.URL.Query().Get(key); v != "" {
|
|
if n, err := strconv.Atoi(v); err == nil && n > 0 {
|
|
return n
|
|
}
|
|
}
|
|
return def
|
|
}
|