refactor: remove legacy multi-agent-mux skills and infrastructure, and add Mattermost notification script and collaboration documentation
This commit is contained in:
@@ -195,6 +195,7 @@ message AgentTask {
|
||||
bytes payload = 5; // Protobuf Any 직렬화
|
||||
map<string, string> metadata = 6; // X-Tenant-ID, X-Device-Class 등
|
||||
google.protobuf.Timestamp deadline = 7;
|
||||
string control_intent_key = 8; // 비멱등 제어 명령의 중복 필터링을 위한 CIK
|
||||
}
|
||||
|
||||
// 에이전트 태스크 응답
|
||||
@@ -422,14 +423,16 @@ func New(ctx context.Context, cfg GatewayConfig) (*Gateway, error) {
|
||||
return nil, fmt.Errorf("SPIFFE credentials: %w", err)
|
||||
}
|
||||
|
||||
// 2. Resume Token 관리자 초기화
|
||||
// 2. Resume Token 및 Idempotency 관리자 초기화
|
||||
tokenMgr := NewResumeTokenManager()
|
||||
idempotencyMgr := NewIdempotencyManager()
|
||||
|
||||
// 3. gRPC 서버 인터셉터 체인 구성
|
||||
srv := grpc.NewServer(
|
||||
grpc.Creds(creds),
|
||||
grpc.ChainUnaryInterceptor(
|
||||
interceptors.DeadlineEnforcer(defaultDeadlines),
|
||||
interceptors.IdempotencyUnary(idempotencyMgr),
|
||||
interceptors.ResumeTokenUnary(tokenMgr),
|
||||
interceptors.OTelUnary(),
|
||||
),
|
||||
@@ -1086,6 +1089,154 @@ type resumeServerStream struct {
|
||||
func (s *resumeServerStream) Context() context.Context { return s.ctx }
|
||||
```
|
||||
|
||||
#### Idempotency gRPC Interceptor (`internal/gateway/interceptors/idempotency.go`)
|
||||
|
||||
```go
|
||||
package interceptors
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/metadata"
|
||||
"google.golang.org/grpc/status"
|
||||
|
||||
agentv1 "github.com/your-org/edge-aiot-mas/gen/go/agent/v1"
|
||||
)
|
||||
|
||||
const controlIntentKeyHeader = "grpc-metadata-control-intent-key"
|
||||
|
||||
// CacheEntry는 CIK 캐시의 레코드를 정의한다.
|
||||
type CacheEntry struct {
|
||||
State string // "RUNNING", "COMPLETED"
|
||||
ResultPayload interface{} // 캐시된 결과 페이로드 (AgentResult)
|
||||
CreatedAt time.Time
|
||||
}
|
||||
|
||||
// IdempotencyManager는 CIK 기반 중복 제거 필터를 총괄한다.
|
||||
type IdempotencyManager struct {
|
||||
mu sync.RWMutex
|
||||
cache map[string]*CacheEntry
|
||||
ttl time.Duration
|
||||
}
|
||||
|
||||
// NewIdempotencyManager는 IdempotencyManager를 초기화한다.
|
||||
func NewIdempotencyManager() *IdempotencyManager {
|
||||
mgr := &IdempotencyManager{
|
||||
cache: make(map[string]*CacheEntry),
|
||||
ttl: 1 * time.Hour, // 기본 TTL 1시간
|
||||
}
|
||||
// 백그라운드에서 캐시 만료 정리 고루틴 기동
|
||||
go mgr.cleanupLoop()
|
||||
return mgr
|
||||
}
|
||||
|
||||
func (m *IdempotencyManager) cleanupLoop() {
|
||||
ticker := time.NewTicker(10 * time.Minute)
|
||||
for range ticker.C {
|
||||
m.mu.Lock()
|
||||
now := time.Now()
|
||||
for k, v := range m.cache {
|
||||
if now.Sub(v.CreatedAt) > m.ttl {
|
||||
delete(m.cache, k)
|
||||
}
|
||||
}
|
||||
m.mu.Unlock()
|
||||
}
|
||||
}
|
||||
|
||||
// Get은 CIK에 해당하는 캐시 데이터를 조회한다.
|
||||
func (m *IdempotencyManager) Get(cik string) (*CacheEntry, bool) {
|
||||
m.mu.RLock()
|
||||
defer m.mu.RUnlock()
|
||||
entry, ok := m.cache[cik]
|
||||
return entry, ok
|
||||
}
|
||||
|
||||
// Set은 CIK에 대한 캐시 레코드를 등록하거나 갱신한다.
|
||||
func (m *IdempotencyManager) Set(cik string, entry *CacheEntry) {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
entry.CreatedAt = time.Now()
|
||||
m.cache[cik] = entry
|
||||
}
|
||||
|
||||
// IdempotencyUnary는 비멱등 Unary 명령의 중복 실행을 방지하는 인터셉터다.
|
||||
func IdempotencyUnary(mgr *IdempotencyManager) grpc.UnaryServerInterceptor {
|
||||
return func(
|
||||
ctx context.Context,
|
||||
req interface{},
|
||||
info *grpc.UnaryServerInfo,
|
||||
handler grpc.UnaryHandler,
|
||||
) (interface{}, error) {
|
||||
// 1. 요청 메시지가 AgentTask인지 타입 단언
|
||||
task, ok := req.(*agentv1.AgentTask)
|
||||
if !ok {
|
||||
return handler(ctx, req)
|
||||
}
|
||||
|
||||
// 2. 비멱등(Non-idempotent) 제어 명령 유형인지 검증 (infer, control, ota 등)
|
||||
if task.TaskType != "control" && task.TaskType != "ota" {
|
||||
return handler(ctx, req)
|
||||
}
|
||||
|
||||
// 3. Control Intent Key (CIK) 추출
|
||||
cik := task.ControlIntentKey
|
||||
if cik == "" {
|
||||
// 들어오는 메타데이터 헤더에서 추출 시도
|
||||
if md, ok := metadata.FromIncomingContext(ctx); ok {
|
||||
keys := md.Get(controlIntentKeyHeader)
|
||||
if len(keys) > 0 {
|
||||
cik = keys[0]
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// CIK가 비어 있으면 사전 검증을 통과시킬 수 없으므로 거부하거나 바이패스
|
||||
if cik == "" {
|
||||
return nil, status.Error(codes.InvalidArgument, "Control Intent Key (CIK) is required for non-idempotent tasks")
|
||||
}
|
||||
|
||||
// 4. 캐시 조회 및 사전 검증(Pre-flight Check)
|
||||
if entry, hit := mgr.Get(cik); hit {
|
||||
switch entry.State {
|
||||
case "RUNNING":
|
||||
return nil, status.Error(codes.Aborted, "Operation is already in progress under this Control Intent Key")
|
||||
case "COMPLETED":
|
||||
// 중복 동작 방지: 캐시된 결과 즉시 반환
|
||||
return entry.ResultPayload, nil
|
||||
}
|
||||
}
|
||||
|
||||
// 5. 캐시에 'RUNNING' 상태로 임시 선점
|
||||
mgr.Set(cik, &CacheEntry{State: "RUNNING"})
|
||||
|
||||
// 6. 핸들러 실행 (하부 물리 계층 명령 전달)
|
||||
resp, err := handler(ctx, req)
|
||||
if err != nil {
|
||||
// 실패 시 캐시 레코드 삭제하여 재시도 허용
|
||||
m := mgr
|
||||
m.mu.Lock()
|
||||
delete(m.cache, cik)
|
||||
m.mu.Unlock()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// 7. 성공 결과 캐싱 완료 처리
|
||||
mgr.Set(cik, &CacheEntry{
|
||||
State: "COMPLETED",
|
||||
ResultPayload: resp,
|
||||
})
|
||||
|
||||
return resp, nil
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
|
||||
### 4.4 A2A HTTP/3 엔드포인트
|
||||
|
||||
#### `internal/gateway/a2a_handler.go`
|
||||
|
||||
Reference in New Issue
Block a user