How to Implement a New Data Source Connector in WeKnora: A Complete Guide Following Feishu, GitLab, and Notion Patterns

To implement a new data source connector in WeKnora, create a package under internal/datasource/connector/<your-source>/, implement the datasource.Connector interface with methods like Type, Validate, ListResources, and FetchAll, wrap the external API in a dedicated client, and register the connector metadata in internal/datasource/connector.go and internal/types/types.go.

WeKnora synchronizes external knowledge bases through a pluggable connector architecture that follows strict interface contracts. Whether you are integrating a custom document management system or a third-party wiki platform, the repository provides established patterns in the Feishu Wiki, GitLab, and Notion connectors that serve as production-ready blueprints for your implementation.

Understanding the Connector Architecture

The Connector Interface Contract

Every data source connector must fulfill the datasource.Connector interface defined in [internal/datasource/connector.go](https://github.com/Tencent/WeKnora/blob/main/internal/datasource/connector.go). This contract ensures uniform behavior across all integrations, enabling the sync engine to discover, validate, and extract content without source-specific logic.

The interface requires six core methods:

  • Type() – Returns a unique string identifier (constant defined in internal/types/types.go)
  • Validate() – Verifies authentication credentials through a lightweight API call
  • ListResources() – Enumerates available collections or documents, supporting hierarchical browsing
  • ResolveResourceAncestors() – Walks up the parent chain to build breadcrumb paths for the resource picker
  • FetchAll() – Performs full synchronization by downloading all specified resources
  • FetchIncremental() – Supports delta sync using cursor-based state comparison

Streaming and Incremental Sync Support

For large datasets that risk memory exhaustion, WeKnora supports the datasource.StreamingConnector extension. This optional interface adds the FetchStream() method, which emits items one-by-one via a StreamHandler callback rather than returning complete slices. The Feishu Wiki connector demonstrates this pattern through its core.FetchStreamEngine integration, while the GitLab connector implements FetchStream directly in [connector.go](https://github.com/Tencent/WeKnora/blob/main/internal/datasource/connector/gitlab/connector.go#L28).

Step-by-Step Implementation Guide

Step 1: Create the Connector Package Structure

Create a new directory under internal/datasource/connector/<your-source>/ containing:

  • connector.go – The main implementation file
  • client.go – Thin wrapper around the external API
  • connector_test.go – Unit tests for validation and fetching logic
  • types.go (optional) – Source-specific structs and cursor definitions

Use the existing feishu/wiki, gitlab, or notion packages as architectural references. The Feishu Wiki connector in [internal/datasource/connector/feishu/wiki/connector.go](https://github.com/Tencent/WeKnora/blob/main/internal/datasource/connector/feishu/wiki/connector.go) exemplifies the standard layout.

Step 2: Implement the Connector Interface Methods

Define a struct (typically holding your API client) and implement all six required methods. The FetchIncremental method must handle types.SyncCursor—an opaque map storing synchronization state between runs.

// internal/datasource/connector/acmedocs/connector.go
package acmedocs

import (
	"context"
	"fmt"

	"github.com/Tencent/WeKnora/internal/datasource"
	"github.com/Tencent/WeKnora/internal/types"
)

type Connector struct {
	client *client
}

func NewConnector() *Connector { return &Connector{} }

func (c *Connector) Type() string { 
	return types.ConnectorTypeAcmeDocs 
}

func (c *Connector) Validate(ctx context.Context, cfg *types.DataSourceConfig) error {
	acmeCfg, err := parseAcmeConfig(cfg)
	if err != nil { return err }
	c.client, err = newClient(acmeCfg.APIKey, acmeCfg.BaseURL)
	if err != nil { return err }
	return c.client.Ping(ctx)
}

func (c *Connector) ListResources(ctx context.Context, cfg *types.DataSourceConfig, parentID string) ([]types.Resource, error) {
	if parentID == "" {
		return c.listCollections(ctx)
	}
	return c.listDocuments(ctx, parentID)
}

func (c *Connector) ResolveResourceAncestors(ctx context.Context, cfg *types.DataSourceConfig, ids []string) ([]string, error) {
	var ancestors []string
	seen := map[string]bool{}
	for _, id := range ids {
		cur := id
		for cur != "" && !seen[cur] {
			ancestors = append(ancestors, cur)
			seen[cur] = true
			parent, err := c.client.GetParent(ctx, cur)
			if err != nil { break }
			cur = parent.ID
		}
	}
	return ancestors, nil
}

func (c *Connector) FetchAll(ctx context.Context, cfg *types.DataSourceConfig, ids []string) ([]types.FetchedItem, error) {
	var result []types.FetchedItem
	for _, colID := range ids {
		docs, err := c.client.ListDocs(ctx, colID)
		if err != nil { return nil, err }
		for _, d := range docs {
			body, _ := c.client.DownloadDoc(ctx, d.ID)
			result = append(result, types.FetchedItem{
				ExternalID: d.ID,
				Title:      d.Title,
				Content:    body,
				ContentType: "text/markdown",
				FileName:   d.Title + ".md",
				URL:        d.URL,
				Metadata: map[string]string{
					"channel": types.ConnectorTypeAcmeDocs,
				},
			})
		}
	}
	return result, nil
}

func (c *Connector) FetchIncremental(ctx context.Context, cfg *types.DataSourceConfig, cur *types.SyncCursor) ([]types.FetchedItem, *types.SyncCursor, error) {
	return nil, nil, fmt.Errorf("incremental sync not yet implemented")
}

Step 3: Build the API Client Layer

The client abstracts authentication, pagination, and error handling. It should expose only the minimal functions your connector requires (e.g., listSpaces, downloadFile). Follow the pattern in [internal/datasource/connector/feishu/core/client.go](https://github.com/Tencent/WeKnora/blob/main/internal/datasource/connector/feishu/core/client.go) or [internal/datasource/connector/gitlab/client.go](https://github.com/Tencent/WeKnora/blob/main/internal/datasource/connector/gitlab/client.go).

// internal/datasource/connector/acmedocs/client.go
package acmedocs

import (
	"context"
	"fmt"
	"net/http"
)

type client struct {
	baseURL string
	apiKey  string
	http    *http.Client
}

func newClient(baseURL, apiKey string) (*client, error) {
	return &client{
		baseURL: baseURL, 
		apiKey:  apiKey, 
		http:    http.DefaultClient,
	}, nil
}

func (c *client) Ping(ctx context.Context) error {
	req, _ := http.NewRequestWithContext(ctx, "GET", c.baseURL+"/ping", nil)
	req.Header.Set("Authorization", "Bearer "+c.apiKey)
	resp, err := c.http.Do(req)
	if err != nil { return err }
	if resp.StatusCode != 200 { 
		return fmt.Errorf("unexpected status %d", resp.StatusCode) 
	}
	return nil
}

Step 4: Add Streaming Support (Optional)

If the source supports it, implement FetchStream to enable memory-bounded synchronization. The method receives a datasource.StreamHandler with Emit() and Checkpoint() methods.

func (c *Connector) FetchStream(ctx context.Context, cfg *types.DataSourceConfig, cur *types.SyncCursor, h datasource.StreamHandler) (*types.SyncCursor, error) {
	collections, _ := c.client.ListCollections(ctx)
	for _, col := range collections {
		docs, _ := c.client.ListDocs(ctx, col.ID)
		for _, d := range docs {
			body, _ := c.client.DownloadDoc(ctx, d.ID)
			h.Emit(types.FetchedItem{
				ExternalID: d.ID,
				Title:      d.Title,
				Content:    body,
			})
		}
		h.Checkpoint(col.ID) // Save progress after each collection
	}
	return cur, nil
}

Step 5: Register the Connector in the Registry

Update internal/datasource/connector.go to add metadata to ConnectorMetadataRegistry:

types.ConnectorTypeAcmeDocs: {
	Type:         types.ConnectorTypeAcmeDocs,
	Name:         "AcmeDocs",
	Description:  "Sync documents from AcmeDocs platform",
	Priority:     9,
	AuthType:     "api_key",
	Capabilities: []string{"incremental"},
},

Then add the constant in internal/types/types.go:

const ConnectorTypeAcmeDocs = "acmedocs"

Step 6: Write Comprehensive Tests

Create connector_test.go covering validation, resource listing, and both full and incremental fetching. Reference the test suite in [internal/datasource/connector/gitlab/connector_test.go](https://github.com/Tencent/WeKnora/blob/main/internal/datasource/connector/gitlab/connector_test.go) for assertion patterns and mock client strategies.

Key Implementation Patterns from Reference Connectors

Feishu Wiki demonstrates sophisticated streaming with FetchStream delegating to core.FetchStreamEngine for concurrent document processing. GitLab provides the canonical example of cursor-based incremental sync using gitLabCursor structs that encode revision IDs. Notion shows lightweight resource ID handling using raw page IDs without complex hierarchy encoding.

Each connector defines its own resource ID encoding strategy:

  • GitLab uses projectID:path notation
  • Feishu uses spaceID:nodeToken pairs
  • Notion uses raw page identifiers

Helper functions like parseWikiResourceID centralize parsing logic and prevent ID format fragmentation.

Summary

  • Interface compliance: Implement all six methods of datasource.Connector defined in internal/datasource/connector.go
  • Client separation: Isolate HTTP logic in a dedicated client file handling auth, retries, and pagination
  • Registration: Add metadata to ConnectorMetadataRegistry and define the type constant in internal/types/types.go
  • Streaming: Optionally implement StreamingConnector for large dataset processing without memory spikes
  • Testing: Mirror the test coverage in existing connectors like GitLab to verify validation, listing, and fetching logic
  • Reference implementation: Study Feishu Wiki for streaming patterns and GitLab for cursor-based incremental sync

Frequently Asked Questions

What methods are required to implement the datasource.Connector interface in WeKnora?

You must implement six methods: Type() to return the connector identifier, Validate() to check credentials, ListResources() to enumerate available items, ResolveResourceAncestors() to build hierarchical paths, FetchAll() for full synchronizations, and FetchIncremental() for delta updates using cursors. These are defined in [internal/datasource/connector.go](https://github.com/Tencent/WeKnora/blob/main/internal/datasource/connector.go).

How does WeKnora handle incremental synchronization in connectors?

Incremental sync uses the types.SyncCursor map to store opaque state between runs. The FetchIncremental method receives the previous cursor, compares it with current source state (typically via timestamps or revision IDs), fetches only changed items, and returns a new cursor encoding the new state. The GitLab connector demonstrates this with its gitLabCursor implementation.

What is the difference between Connector and StreamingConnector in WeKnora?

Connector requires methods that return complete slices of results, which can exhaust memory for large datasets. StreamingConnector extends this with FetchStream(), which emits items individually through a StreamHandler callback and supports checkpointing progress, enabling constant-memory synchronization of large repositories.

Where do I register a new connector type in the WeKnora codebase?

Registration requires two updates: add the connector type constant (e.g., ConnectorTypeAcmeDocs) to internal/types/types.go, and append the metadata struct to the ConnectorMetadataRegistry map in internal/datasource/connector.go. This makes the connector discoverable by the sync engine and administrative interfaces.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →