package shunyunbao import ( "context" "encoding/json" "reflect" "strconv" "cmroubao/backend-api/internal/domain" ) const ( sourcePageSize = 20 sourceMaxMatches = 100 ) // QueryOrder implements the freight source port using the manager's single // in-memory ERP session. It never exposes the ERP list or detail responses. func (manager *SessionManager) QueryOrder( ctx context.Context, orderNumber string, ) (domain.FreightSourceBatch, error) { query, err := OrderNumberQuery(orderNumber) if err != nil { return domain.FreightSourceBatch{}, domain.ErrFreightSourceProtocol } records, err := manager.queryRecords(ctx, query) if err != nil { return domain.FreightSourceBatch{}, err } if len(records) == 0 { return domain.FreightSourceBatch{}, domain.ErrFreightSourceNotFound } return NormalizeFreightResult( RawQuery{Mode: domain.FreightSyncOrderNumber}, records, ) } // QueryCreatedRange implements the bounded Asia/Shanghai date-window source // port. The caller owns larger-window splitting and watermark semantics. func (manager *SessionManager) QueryCreatedRange( ctx context.Context, createdFrom, createdTo string, ) (domain.FreightSourceBatch, error) { query, err := CreatedRangeQuery(createdFrom, createdTo) if err != nil { return domain.FreightSourceBatch{}, domain.ErrFreightSourceProtocol } records, err := manager.queryRecords(ctx, query) if err != nil { return domain.FreightSourceBatch{}, err } return NormalizeFreightResult( RawQuery{ Mode: domain.FreightSyncCreatedRange, CreatedFrom: createdFrom, CreatedTo: createdTo, }, records, ) } func (manager *SessionManager) queryRecords( ctx context.Context, query QueryCondition, ) ([]RawRecord, error) { manager.mu.Lock() defer manager.mu.Unlock() if !manager.configuredLocked() { return nil, domain.ErrFreightSourceNotConfigured } if !manager.authenticated { return nil, domain.ErrFreightSourceSessionNeeded } if _, err := manager.validateLocked(ctx); err != nil { return nil, err } stocks, err := manager.queryStocksLocked(ctx, query) if err != nil { return nil, err } if len(stocks) == 0 { return []RawRecord{}, nil } details, err := manager.queryDetailsLocked(ctx, stocks) if err != nil { return nil, err } records := make([]RawRecord, 0, len(stocks)) for _, stock := range stocks { stockID, err := sourceExternalID(stock["id"]) if err != nil { return nil, domain.ErrFreightSourceProtocol } detail, exists := details[stockID] if !exists { return nil, domain.ErrFreightSourceProtocol } records = append(records, RawRecord{Stock: stock, Detail: detail}) } return records, nil } func (manager *SessionManager) queryStocksLocked( ctx context.Context, query QueryCondition, ) ([]map[string]any, error) { firstPayload, err := StockListPayload(query, 0, 1, sourcePageSize) if err != nil { return nil, domain.ErrFreightSourceProtocol } firstBody, err := json.Marshal(firstPayload) if err != nil { return nil, domain.ErrFreightSourceProtocol } totalData, err := manager.requestJSONLocked( ctx, "POST", StockListTotalPath, firstBody, false, ) if err != nil { return nil, err } total, err := sourceTotal(totalData) if err != nil || total > sourceMaxMatches { return nil, domain.ErrFreightSourceProtocol } if total == 0 { return []map[string]any{}, nil } rows := make([]map[string]any, 0, total) for start := 0; start < total; start += sourcePageSize { payload, err := StockListPayload( query, start, start/sourcePageSize+1, sourcePageSize, ) if err != nil { return nil, domain.ErrFreightSourceProtocol } body, err := json.Marshal(payload) if err != nil { return nil, domain.ErrFreightSourceProtocol } data, err := manager.requestJSONLocked( ctx, "POST", StockListPath, body, false, ) if err != nil { return nil, err } page, err := sourceObjectList(data) if err != nil || len(page) > sourcePageSize { return nil, domain.ErrFreightSourceProtocol } rows = append(rows, page...) } if len(rows) != total { return nil, domain.ErrFreightSourceProtocol } stocksByID := make(map[string]map[string]any, len(rows)) stocks := make([]map[string]any, 0, len(rows)) for _, stock := range rows { stockID, err := sourceExternalID(stock["id"]) if err != nil { return nil, domain.ErrFreightSourceProtocol } if existing, exists := stocksByID[stockID]; exists { if !reflect.DeepEqual(existing, stock) { return nil, domain.ErrFreightSourceProtocol } continue } stocksByID[stockID] = stock stocks = append(stocks, stock) } return stocks, nil } func (manager *SessionManager) queryDetailsLocked( ctx context.Context, stocks []map[string]any, ) (map[string]map[string]any, error) { ids := make([]string, 0, len(stocks)) for _, stock := range stocks { stockID, err := sourceExternalID(stock["id"]) if err != nil { return nil, domain.ErrFreightSourceProtocol } ids = append(ids, stockID) } detailsByID := make(map[string]map[string]any, len(ids)) requestedIDs := make(map[string]struct{}, len(ids)) for _, id := range ids { requestedIDs[id] = struct{}{} } for start := 0; start < len(ids); start += maxDetailIDs { end := start + maxDetailIDs if end > len(ids) { end = len(ids) } payload, err := StockDetailPayload(ids[start:end]) if err != nil { return nil, domain.ErrFreightSourceProtocol } body, err := json.Marshal(payload) if err != nil { return nil, domain.ErrFreightSourceProtocol } data, err := manager.requestJSONLocked( ctx, "POST", StockDetailPath+"?hist=0", body, false, ) if err != nil { return nil, err } values, err := sourceObjectList(data) if err != nil { return nil, domain.ErrFreightSourceProtocol } for _, detail := range values { detailID, err := sourceExternalID(detail["id"]) if err != nil { return nil, domain.ErrFreightSourceProtocol } if _, requested := requestedIDs[detailID]; !requested { return nil, domain.ErrFreightSourceProtocol } if existing, exists := detailsByID[detailID]; exists && !reflect.DeepEqual(existing, detail) { return nil, domain.ErrFreightSourceProtocol } detailsByID[detailID] = detail } } if len(detailsByID) != len(ids) { return nil, domain.ErrFreightSourceProtocol } return detailsByID, nil } func sourceTotal(value any) (int, error) { var text string switch typed := value.(type) { case json.Number: text = string(typed) case string: text = typed default: return 0, strconv.ErrSyntax } parsed, err := strconv.ParseInt(text, 10, 0) if err != nil || parsed < 0 || parsed > int64(sourceMaxMatches) { return 0, strconv.ErrSyntax } return int(parsed), nil } func sourceObjectList(value any) ([]map[string]any, error) { data, ok := value.(map[string]any) if !ok { return nil, strconv.ErrSyntax } rawList, ok := data["list"].([]any) if !ok { return nil, strconv.ErrSyntax } result := make([]map[string]any, 0, len(rawList)) for _, value := range rawList { row, ok := value.(map[string]any) if !ok { return nil, strconv.ErrSyntax } result = append(result, row) } return result, nil } func sourceExternalID(value any) (string, error) { return externalID(value) }