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 ininternal/types/types.go)Validate()– Verifies authentication credentials through a lightweight API callListResources()– Enumerates available collections or documents, supporting hierarchical browsingResolveResourceAncestors()– Walks up the parent chain to build breadcrumb paths for the resource pickerFetchAll()– Performs full synchronization by downloading all specified resourcesFetchIncremental()– 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 fileclient.go– Thin wrapper around the external APIconnector_test.go– Unit tests for validation and fetching logictypes.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:pathnotation - Feishu uses
spaceID:nodeTokenpairs - 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.Connectordefined ininternal/datasource/connector.go - Client separation: Isolate HTTP logic in a dedicated client file handling auth, retries, and pagination
- Registration: Add metadata to
ConnectorMetadataRegistryand define the type constant ininternal/types/types.go - Streaming: Optionally implement
StreamingConnectorfor 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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →