Creating Plugins
Extend Xberg with custom extractors, post-processors, OCR backends, and validators registered globally for use across all extraction calls.
Plugin Types
Section titled “Plugin Types”| Type | Purpose | Use case |
|---|---|---|
| DocumentExtractor | Extract content from file formats | New format support, override built-in extractors |
| PostProcessor | Transform extraction results | Metadata enrichment, content filtering, text normalization |
| OcrBackend | Perform OCR on images | Cloud OCR services, custom OCR engines |
| Validator | Validate extraction quality | Minimum content length, quality score thresholds |
| EmbeddingBackend | Generate embedding vectors | Custom embedding models, RAG pipelines |
| RerankerBackend | Score query/document pairs | Cross-encoder reranking of retrieved chunks |
| TokenizerBackend | Count tokens for chunk boundaries | Model-specific tokenizers, token-aware chunk sizing |
| Renderer | Convert results to output formats | Custom Markdown, HTML, Djot, or plain-text renderers |
All plugins must be thread-safe (Send + Sync in Rust, thread-safe in Python) and
implement initialize() / shutdown() lifecycle methods.
Document Extractors
Section titled “Document Extractors”Implementation
Section titled “Implementation”use xberg::plugins::{DocumentExtractor, Plugin};use xberg::{Result, ExtractedDocument, ExtractInput, ExtractionConfig};use async_trait::async_trait;
struct CustomJsonExtractor;
impl Plugin for CustomJsonExtractor { fn name(&self) -> &str { "custom-json-extractor" } fn version(&self) -> String { "1.0.0".to_string() } fn initialize(&self) -> Result<()> { Ok(()) } fn shutdown(&self) -> Result<()> { Ok(()) }}
#[async_trait]impl DocumentExtractor for CustomJsonExtractor { async fn extract( &self, input: ExtractInput, _config: &ExtractionConfig, ) -> Result<ExtractedDocument> { let bytes = input.bytes.unwrap_or_default(); let json: serde_json::Value = serde_json::from_slice(&bytes)?; let text = extract_text_from_json(&json);
let mut document = ExtractedDocument::default(); document.content = text; document.mime_type = "application/json".into(); Ok(document) }
fn supported_mime_types(&self) -> &[&str] { &["application/json", "text/json"] }
fn priority(&self) -> i32 { 50 }}
fn extract_text_from_json(value: &serde_json::Value) -> String { match value { serde_json::Value::String(s) => format!("{}\n", s), serde_json::Value::Array(arr) => arr.iter().map(extract_text_from_json).collect(), serde_json::Value::Object(obj) => obj.values().map(extract_text_from_json).collect(), _ => String::new(), }}from xberg import register_document_extractor, ExtractInput, ExtractionConfigimport json
class CustomJsonExtractor: def name(self) -> str: return "custom-json-extractor"
def version(self) -> str: return "1.0.0"
def supported_mime_types(self) -> list[str]: return ["application/json"]
def priority(self) -> int: return 50
def extract(self, input: ExtractInput, config: ExtractionConfig) -> dict: data: dict = json.loads(input.bytes) text: str = self._extract_text(data) return {"content": text, "mime_type": "application/json"}
def _extract_text(self, obj: object) -> str: if isinstance(obj, str): return f"{obj}\n" if isinstance(obj, list): return "".join(self._extract_text(item) for item in obj) if isinstance(obj, dict): return "".join(self._extract_text(v) for v in obj.values()) return ""
def initialize(self) -> None: pass
def shutdown(self) -> None: pass
extractor: CustomJsonExtractor = CustomJsonExtractor()register_document_extractor(extractor)Registration
Section titled “Registration”from xberg import register_document_extractor, ExtractInput, ExtractionConfig, ExtractedDocument
class CustomExtractor: def name(self) -> str: return "custom"
def version(self) -> str: return "1.0.0"
def supported_mime_types(self) -> list[str]: return ["application/x-custom"]
def extract(self, input: ExtractInput, config: ExtractionConfig) -> dict: content = input.bytes.decode("utf-8") if input.bytes else "" return {"content": content, "mime_type": "application/x-custom"}
extractor = CustomExtractor()register_document_extractor(extractor)print("Extractor registered")import { listDocumentExtractors, registerDocumentExtractor, unregisterDocumentExtractor, clearDocumentExtractors, type DocumentExtractor, type ExtractedDocument,} from "@xberg-io/xberg";
// Custom document extractors are supported: implement `DocumentExtractor`// and register it. See `plugin_extractor.md` for a complete example.const customExtractor: DocumentExtractor = { name: () => "custom-text-extractor", supportedMimeTypes: () => ["text/x-custom"], priority: () => 60, async extract(): Promise<ExtractedDocument> { return { content: "custom extraction result", mimeType: "text/x-custom" }; },};registerDocumentExtractor(customExtractor);
// List all registered document extractorsconst extractors = listDocumentExtractors();console.log("Available extractors:", extractors);
// Unregister a specific extractor (use with caution)unregisterDocumentExtractor("custom-text-extractor");
// Clear all extractors (use with extreme caution)// clearDocumentExtractors();use xberg::plugins::{Plugin, DocumentExtractor};use xberg::{ExtractInput, ExtractionConfig, ExtractedDocument, Result, register_document_extractor};use async_trait::async_trait;use std::sync::Arc;
struct CustomJsonExtractor;
impl Plugin for CustomJsonExtractor { fn name(&self) -> &str { "custom-json-extractor" } fn version(&self) -> String { "1.0.0".to_string() } fn initialize(&self) -> Result<()> { Ok(()) } fn shutdown(&self) -> Result<()> { Ok(()) }}
#[async_trait]impl DocumentExtractor for CustomJsonExtractor { async fn extract(&self, input: ExtractInput, _config: &ExtractionConfig) -> Result<ExtractedDocument> { let bytes = input.bytes.unwrap_or_default(); let mut document = ExtractedDocument::default(); document.content = String::from_utf8_lossy(&bytes).to_string(); document.mime_type = "application/json".into(); Ok(document) }
fn supported_mime_types(&self) -> &[&str] { &["application/json", "text/json"] }}
fn register_custom_extractor() -> Result<()> { let extractor = Arc::new(CustomJsonExtractor); register_document_extractor(extractor)?; Ok(())}package main
import ( "log"
"github.com/xberg-io/xberg/packages/go")
// jsonExtractor is a minimal custom xberg.DocumentExtractor for JSON documents.type jsonExtractor struct{}
func (jsonExtractor) Name() string { return "custom-json-extractor" }func (jsonExtractor) Version() string { return "1.0.0" }func (jsonExtractor) Initialize() error { return nil }func (jsonExtractor) Shutdown() error { return nil }func (jsonExtractor) Priority() int32 { return 50 }
func (jsonExtractor) CanHandle(_path string, mimeType string) bool { return mimeType == "application/json"}
func (jsonExtractor) Extract(input xberg.ExtractInput, config xberg.ExtractionConfig) (xberg.ExtractedDocument, error) { return xberg.ExtractedDocument{}, nil}
func (jsonExtractor) SupportedMimeTypes() []string { return []string{"application/json"}}
func main() { // Register custom extractor if err := xberg.RegisterDocumentExtractor(jsonExtractor{}); err != nil { log.Fatalf("register extractor failed: %v", err) }
input := xberg.ExtractInputFromURI("document.json") result, err := xberg.Extract(*input, xberg.ExtractionConfig{}) if err != nil { log.Fatalf("extract failed: %v", err) } log.Printf("Extracted content length: %d", len(result.Results[0].Content))}import io.xberg.Xberg;import io.xberg.ExtractInputKind;import io.xberg.ExtractionResult;import io.xberg.ExtractedDocument;import io.xberg.ExtractInput;import io.xberg.ExtractionConfig;import io.xberg.XbergRsException;
public class CustomExtractorExample { public static void main(String[] args) { try { ExtractionResult output = Xberg.extract( ExtractInput.builder().withKind(ExtractInputKind.URI).withUri("document.json").build(), ExtractionConfig.builder().build() ); ExtractedDocument result = output.results().get(0); System.out.println("Extracted content length: " + result.content().length()); } catch (XbergRsException e) { e.printStackTrace(); } }}using Xberg;using System;using System.Collections.Generic;
var extractor = new CustomExtractor();DocumentExtractorRegistry.RegisterDocumentExtractor(extractor);Console.WriteLine("Extractor registered");
public class CustomExtractor : IDocumentExtractor{ public string Name => "custom"; public string Version => "1.0.0"; public int Priority => 50; public List<string> SupportedMimeTypes => new() { "application/x-custom" };
public void Initialize() { } public void Shutdown() { }
public bool CanHandle(string path, string mimeType) => mimeType == "application/x-custom";
public ExtractedDocument Extract(ExtractInput input, ExtractionConfig config) { return new ExtractedDocument { Content = "Extracted content", MimeType = "application/x-custom", Metadata = new Metadata(), }; }}require 'xberg'
# Register custom extractor with priority 50Xberg.register_document_extractor( name: "custom-json-extractor", extractor: ->(content, mime_type, config) { JSON.parse(content.to_s) }, priority: 50)
input = Xberg::ExtractInput.new(uri: "document.json")config = Xberg::ExtractionConfig.newresult = Xberg.extract(input, config)puts "Extracted content length: #{result.results.first.content.length}"Priority System
Section titled “Priority System”When multiple extractors support the same MIME type, the highest priority wins:
| Range | Level |
|---|---|
| 0–25 | Fallback / low-quality |
| 26–49 | Alternative |
| 50 | Default (built-in) |
| 51–75 | Enhanced / premium |
| 76–100 | Specialized / high-priority |
Post-Processors
Section titled “Post-Processors”Processors execute in three stages:
- Early — Foundational: language detection, quality scoring, text normalization
- Middle — Transformation: keyword extraction, token reduction, summarization
- Late — Final: custom metadata, analytics, output formatting
Implementation
Section titled “Implementation”use xberg::plugins::{Plugin, PostProcessor, ProcessingStage};use xberg::{Result, ExtractedDocument, ExtractionConfig, ProcessingWarning};use async_trait::async_trait;
struct WordCountProcessor;
impl Plugin for WordCountProcessor { fn name(&self) -> &str { "word-count" } fn version(&self) -> String { "1.0.0".to_string() } fn initialize(&self) -> Result<()> { Ok(()) } fn shutdown(&self) -> Result<()> { Ok(()) }}
#[async_trait]impl PostProcessor for WordCountProcessor { async fn process( &self, result: &mut ExtractedDocument, _config: &ExtractionConfig ) -> Result<()> { let word_count = result.content.split_whitespace().count();
result.processing_warnings.push(ProcessingWarning { source: "word-count".into(), message: format!("Processed with word count: {}", word_count).into() });
Ok(()) }
fn processing_stage(&self) -> ProcessingStage { ProcessingStage::Early }
fn should_process( &self, result: &ExtractedDocument, _config: &ExtractionConfig ) -> bool { !result.content.is_empty() }}import loggingfrom xberg import register_post_processor, ExtractedDocument, ExtractionConfig
logger = logging.getLogger(__name__)
class WordCountProcessor: def name(self) -> str: return "word_count"
def version(self) -> str: return "1.0.0"
def processing_stage(self) -> str: return "early"
def process(self, result: ExtractedDocument, config: ExtractionConfig) -> None: word_count: int = len(result.content.split()) logger.info(f"Word count: {word_count}")
def should_process(self, result: ExtractedDocument, config: ExtractionConfig) -> bool: return bool(result.content)
def initialize(self) -> None: pass
def shutdown(self) -> None: pass
processor: WordCountProcessor = WordCountProcessor()register_post_processor(processor)alias Xberg.Plugin
# Word Count Post-Processor Plugin# This post-processor automatically counts words in extracted content# and adds the word count to the metadata.
defmodule MyApp.Plugins.WordCountProcessor do @behaviour Xberg.Plugin.PostProcessor require Logger
@impl true def name do "WordCountProcessor" end
@impl true def processing_stage do :post end
@impl true def version do "1.0.0" end
@impl true def initialize do :ok end
@impl true def shutdown do :ok end
@impl true def process(result, _options) do content = result["content"] || "" word_count = content |> String.split(~r/\s+/, trim: true) |> length()
# Update metadata with word count metadata = Map.get(result, "metadata", %{}) updated_metadata = Map.put(metadata, "word_count", word_count)
{:ok, Map.put(result, "metadata", updated_metadata)} endend
# Register the word count post-processorPlugin.register_post_processor(:word_count_processor, MyApp.Plugins.WordCountProcessor)
# Example usageresult = %{"content" => "The quick brown fox jumps over the lazy dog. This is a sample document with multiple words.","metadata" => %{"source" => "document.pdf","pages" => 1}}
case MyApp.Plugins.WordCountProcessor.process(result, %{}) do {:ok, processed_result} -> word_count = processed_result["metadata"]["word_count"] IO.puts("Word count added: #{word_count} words") IO.inspect(processed_result, label: "Processed Result")
{:error, reason} -> IO.puts("Processing failed: #{reason}")end
# List all registered post-processors{:ok, processors} = Plugin.list_post_processors()IO.inspect(processors, label: "Registered Post-Processors")Conditional Processing
Section titled “Conditional Processing”from xberg import ExtractedDocument, ExtractionConfig, register_post_processor
class PdfOnlyProcessor: def name(self) -> str: return "pdf-only-processor"
def version(self) -> str: return "1.0.0"
def processing_stage(self) -> str: return "early"
def process(self, result: ExtractedDocument, config: ExtractionConfig) -> None: pass
def should_process(self, result: ExtractedDocument, config: ExtractionConfig) -> bool: return result.mime_type == "application/pdf"
processor: PdfOnlyProcessor = PdfOnlyProcessor()register_post_processor(processor)use xberg::plugins::{Plugin, PostProcessor, ProcessingStage};use xberg::{ExtractedDocument, ExtractionConfig, Result};use async_trait::async_trait;
struct PdfOnlyProcessor;
impl Plugin for PdfOnlyProcessor { fn name(&self) -> &str { "pdf-only" } fn version(&self) -> String { "1.0.0".to_string() } fn initialize(&self) -> Result<()> { Ok(()) } fn shutdown(&self) -> Result<()> { Ok(()) }}
#[async_trait]impl PostProcessor for PdfOnlyProcessor { async fn process( &self, result: &mut ExtractedDocument, _config: &ExtractionConfig ) -> Result<()> { Ok(()) }
fn processing_stage(&self) -> ProcessingStage { ProcessingStage::Middle }
fn should_process( &self, result: &ExtractedDocument, _config: &ExtractionConfig ) -> bool { result.mime_type == "application/pdf" }}package main
import ( "log"
"github.com/xberg-io/xberg/packages/go")
type pdfOnlyProcessor struct{}
func (processor *pdfOnlyProcessor) Name() string { return "pdf_only_processor" }func (processor *pdfOnlyProcessor) Version() string { return "1.0.0" }func (processor *pdfOnlyProcessor) Initialize() error { return nil }func (processor *pdfOnlyProcessor) Shutdown() error { return nil }func (processor *pdfOnlyProcessor) Priority() int32 { return 70 }func (processor *pdfOnlyProcessor) ProcessingStage() xberg.ProcessingStage { return xberg.ProcessingStageMiddle}func (processor *pdfOnlyProcessor) ShouldProcess( result xberg.ExtractedDocument, _ xberg.ExtractionConfig,) bool { return result.MimeType == "application/pdf"}func (processor *pdfOnlyProcessor) EstimatedDurationMs(_ xberg.ExtractedDocument) uint64 { return 1}func (processor *pdfOnlyProcessor) Process( result xberg.ExtractedDocument, _ xberg.ExtractionConfig,) error { log.Printf("Processing PDF with %d tables", len(result.Tables)) return nil}
func main() { processor := &pdfOnlyProcessor{} if err := xberg.RegisterPostProcessor(processor); err != nil { log.Fatalf("register post-processor: %v", err) } defer func() { if err := xberg.UnregisterPostProcessor(processor.Name()); err != nil { log.Printf("unregister post-processor: %v", err) } }()
for _, path := range []string{"document.pdf", "image.jpg", "spreadsheet.xlsx"} { input := xberg.ExtractInputFromURI(path) result, err := xberg.Extract(*input, xberg.ExtractionConfig{}) if err != nil { log.Printf("extract %s: %v", path, err) continue } log.Printf("Extracted %s as %s", path, result.Results[0].MimeType) }}import io.xberg.ExtractedDocument;import io.xberg.ExtractionConfig;import io.xberg.IPostProcessor;
// Post-processors observe the extracted document (process() returns void);// ExtractedDocument is an immutable record, so this hook can no longer// inject new metadata fields the way older PostProcessor implementations could.IPostProcessor pdfOnly = new IPostProcessor() { @Override public String name() { return "pdf-only"; }
@Override public String version() { return "1.0.0"; }
@Override public void process(ExtractedDocument result, ExtractionConfig config) throws Exception { if (!result.mimeType().equals("application/pdf")) { return; } // Handle PDF-specific processing here. }
@Override public String processing_stage() throws Exception { return "pdf-only"; }
@Override public boolean should_process(ExtractedDocument _result, ExtractionConfig _config) throws Exception { return _result.mimeType().equals("application/pdf"); }
@Override public long estimated_duration_ms(ExtractedDocument _result) throws Exception { return 0; }
@Override public int priority() throws Exception { return 50; }};using Xberg;
public class PdfOnlyProcessor : IPostProcessor{ public string Name => "pdf-only-processor"; public string Version => "1.0.0"; public int Priority => 50; public ProcessingStage ProcessingStage => ProcessingStage.Middle;
public void Initialize() { } public void Shutdown() { }
public ulong EstimatedDurationMs(ExtractedDocument result) => 1;
public void Process(ExtractedDocument result, ExtractionConfig config) { }
public bool ShouldProcess(ExtractedDocument result, ExtractionConfig config) => result.MimeType == "application/pdf";}
class Program{ static void Main() { var processor = new PdfOnlyProcessor(); PostProcessorRegistry.RegisterPostProcessor(processor); }}OCR Backends
Section titled “OCR Backends”Implementation
Section titled “Implementation”use xberg::plugins::{Plugin, OcrBackend, OcrBackendType};use xberg::{Result, ExtractedDocument, OcrConfig};use async_trait::async_trait;
struct CloudOcrBackend { api_key: String, supported_langs: Vec<String>,}
impl Plugin for CloudOcrBackend { fn name(&self) -> &str { "cloud-ocr" } fn version(&self) -> String { "1.0.0".to_string() } fn initialize(&self) -> Result<()> { Ok(()) } fn shutdown(&self) -> Result<()> { Ok(()) }}
#[async_trait]impl OcrBackend for CloudOcrBackend { async fn process_image( &self, image_bytes: &[u8], config: &OcrConfig, ) -> Result<ExtractedDocument> { let language = config.language.first().map(String::as_str).unwrap_or("eng"); let text = self.call_cloud_api(image_bytes, language).await?;
// `ExtractedDocument` has private internal fields, so a struct literal with // `..Default::default()` does not compile outside the crate. Build a default // and assign the public fields instead. let mut document = ExtractedDocument::default(); document.content = text; document.mime_type = "text/plain".into(); Ok(document) }
fn supports_language(&self, lang: &str) -> bool { self.supported_langs.iter().any(|l| l == lang) }
fn backend_type(&self) -> OcrBackendType { OcrBackendType::Custom }
fn supported_languages(&self) -> Vec<String> { self.supported_langs.clone() }}
impl CloudOcrBackend { async fn call_cloud_api( &self, image: &[u8], language: &str ) -> Result<String> { Ok("Extracted text".to_string()) }}from xberg import register_ocr_backend, ExtractedDocument, OcrBackendType, OcrConfig, Metadataimport httpx
class CloudOcrBackend: def __init__(self, api_key: str): self.api_key: str = api_key self.langs: list[str] = ["eng", "deu", "fra"]
def name(self) -> str: return "cloud-ocr"
def version(self) -> str: return "1.0.0"
def supported_languages(self) -> list[str]: return self.langs
def supports_language(self, lang: str) -> bool: return lang in self.langs
def backend_type(self) -> OcrBackendType: return OcrBackendType.CUSTOM
def supports_table_detection(self) -> bool: return False
def supports_document_processing(self) -> bool: return False
def emits_structured_markdown(self) -> bool: return False
def process_image(self, image_bytes: bytes, config: OcrConfig) -> ExtractedDocument: with httpx.Client() as client: response = client.post( "https://api.example.com/ocr", files={"image": image_bytes}, json={"language": config.language[0] if config.language else "eng"}, ) text: str = response.json()["text"] return ExtractedDocument( content=text, mime_type="text/plain", metadata=Metadata(), )
def process_image_file(self, path: str, config: OcrConfig) -> ExtractedDocument: with open(path, "rb") as f: return self.process_image(f.read(), config)
def process_document(self, path: str, config: OcrConfig) -> ExtractedDocument: return self.process_image_file(path, config)
def initialize(self) -> None: pass
def shutdown(self) -> None: pass
backend: CloudOcrBackend = CloudOcrBackend(api_key="your-api-key")register_ocr_backend(backend)import io.xberg.*;import io.xberg.ExtractInputKind;import java.net.http.*;import java.net.URI;import java.nio.file.Files;import java.nio.file.Path;import java.util.List;
public class CloudOcrExample implements IOcrBackend { private final String apiKey;
public CloudOcrExample(String apiKey) { this.apiKey = apiKey; }
@Override public String name() { return "cloud-ocr"; }
@Override public String version() { return "1.0.0"; }
@Override public ExtractedDocument process_image(byte[] image_bytes, OcrConfig config) throws Exception { // Call cloud OCR API HttpClient client = HttpClient.newHttpClient(); HttpRequest request = HttpRequest.newBuilder() .uri(URI.create("https://api.example.com/ocr")) .header("Authorization", "Bearer " + apiKey) .POST(HttpRequest.BodyPublishers.ofByteArray(image_bytes)) .build(); HttpResponse<String> response = client.send(request, HttpResponse.BodyHandlers.ofString()); String text = parseTextFromResponse(response.body()); return ExtractedDocument.builder() .withContent(text) .withMimeType("text/plain") .withMetadata(Metadata.builder().build()) .build(); }
@Override public ExtractedDocument process_image_file(Path path, OcrConfig config) throws Exception { return process_image(Files.readAllBytes(path), config); }
@Override public boolean supports_language(String lang) throws Exception { return true; }
@Override public String backend_type() throws Exception { return "cloud-ocr"; }
@Override public List<String> supported_languages() throws Exception { return List.of("en"); }
// The cloud service offers the same languages whatever the extraction config asks for, so // both config-aware methods delegate to the config-independent ones above. A backend whose // language support depends on the config would read it here instead. @Override public List<String> supported_languages_for(OcrConfig config) throws Exception { return supported_languages(); }
@Override public boolean supports_language_for(OcrConfig config, String language) throws Exception { return supports_language(language); }
@Override public boolean supports_table_detection() throws Exception { return false; }
@Override public boolean supports_document_processing() throws Exception { return false; }
@Override public boolean emits_structured_markdown() throws Exception { return false; }
// The cloud service reports no calibrated page confidence and needs an upright // raster, so declare the least-capable option for both descriptors. @Override public ConfidenceSemantics confidence_semantics() throws Exception { return ConfidenceSemantics.UNCALIBRATED; }
@Override public PageOrientationHandling page_orientation_handling() throws Exception { return PageOrientationHandling.REQUIRES_UPRIGHT; }
@Override public ExtractedDocument process_document(Path _path, OcrConfig _config) throws Exception { throw new UnsupportedOperationException("cloud-ocr does not support whole-document processing"); }
private static String parseTextFromResponse(String json) { // Parse JSON response and extract text field return json; // Simplified }
public static void main(String[] args) { try { Xberg.registerOcrBackend(new CloudOcrExample("your-api-key")); // Use custom OCR backend in extraction // Note: Requires ExtractionConfig with OCR enabled var resultOutput = Xberg.extract( io.xberg.ExtractInput.builder() .withKind(io.xberg.ExtractInputKind.URI) .withUri("scanned.pdf") .build(), io.xberg.ExtractionConfig.builder().build() ); ExtractedDocument result = resultOutput.results().get(0); } catch (Exception e) { e.printStackTrace(); } }}using Xberg;using System;using System.Collections.Generic;using System.Net.Http;using System.Text.Json;using System.Threading.Tasks;
public class CloudOcrBackend : IOcrBackend{ private readonly string _apiKey; private readonly List<string> _langs = new() { "eng", "deu", "fra" }; private readonly HttpClient _httpClient = new();
public CloudOcrBackend(string apiKey) { _apiKey = apiKey; }
public string Name => "cloud-ocr"; public string Version => "1.0.0"; public OcrBackendType BackendType => OcrBackendType.Custom; public List<string> SupportedLanguages => _langs; public bool SupportsTableDetection => false; public bool SupportsDocumentProcessing => false; public bool EmitsStructuredMarkdown => false; // The cloud service reports no calibrated page confidence and needs an upright // raster, so declare the least-capable option for both descriptors. public ConfidenceSemantics ConfidenceSemantics => Xberg.ConfidenceSemantics.Uncalibrated; public PageOrientationHandling PageOrientationHandling => Xberg.PageOrientationHandling.RequiresUpright;
public void Initialize() { } public void Shutdown() => _httpClient.Dispose();
public bool SupportsLanguage(string lang) => _langs.Contains(lang);
// The cloud service offers the same languages whatever the extraction config asks for, so // both config-aware members delegate to the config-independent set above. A backend whose // language support depends on the config (a model chosen per request, say) would read it here. public List<string> SupportedLanguagesFor(OcrConfig config) => _langs;
public bool SupportsLanguageFor(OcrConfig config, string language) => _langs.Contains(language);
public ExtractedDocument ProcessImage(byte[] imageBytes, OcrConfig config) { using var form = new MultipartFormDataContent(); form.Add(new ByteArrayContent(imageBytes), "image"); var lang = config.Language.Count > 0 ? config.Language[0] : "eng"; form.Add(new StringContent(lang), "language");
var response = _httpClient.PostAsync("https://api.example.com/ocr", form).Result; var json = response.Content.ReadAsStringAsync().Result; var doc = JsonDocument.Parse(json); var text = doc.RootElement.GetProperty("text").GetString() ?? "";
return new ExtractedDocument { Content = text, MimeType = "text/plain", Metadata = new Metadata(), }; }
public ExtractedDocument ProcessImageFile(string path, OcrConfig config) => ProcessImage(System.IO.File.ReadAllBytes(path), config);
public ExtractedDocument ProcessDocument(string path, OcrConfig config) => throw new OcrException("cloud-ocr does not support whole-document processing");}
class Program{ static void Main() { var backend = new CloudOcrBackend(apiKey: "your-api-key"); OcrBackendRegistry.RegisterOcrBackend(backend); }}require 'xberg'require 'net/http'
class CloudOcrBackend def name 'cloud-ocr' end
def supported_languages %w[eng fra deu] end
def process_image(image_data, language) uri = URI('https://api.example.com/ocr') req = Net::HTTP::Post.new(uri) req['Authorization'] = "Bearer #{ENV['OCR_API_KEY']}" req.body = image_data res = Net::HTTP.start(uri.hostname, uri.port, use_ssl: true) { |h| h.request(req) } raise StandardError, res.message unless res.is_a?(Net::HTTPSuccess) { content: JSON.parse(res.body)['text'] } rescue StandardError => e raise StandardError, e.message endend
Xberg.register_ocr_backend(CloudOcrBackend.new)config = Xberg::ExtractionConfig.new( ocr: Xberg::OcrConfig.new(backend: 'cloud-ocr'))input = Xberg::ExtractInput.new(uri: 'doc.pdf')Xberg.extract(input, config)Registration
Section titled “Registration”Register the backend and set its name in OcrConfig:
from xberg import register_ocr_backend, unregister_ocr_backend
backend = CloudOcrBackend(api_key="your-api-key")register_ocr_backend(backend)
from xberg import extract, ExtractionConfig, OcrConfig
config = ExtractionConfig(ocr=OcrConfig(backend="cloud-ocr", language=["eng"]))result = extract("scanned.pdf", config=config)
unregister_ocr_backend("cloud-ocr")Validators
Section titled “Validators”use xberg::plugins::{Plugin, Validator};use xberg::{Result, ExtractedDocument, ExtractionConfig, XbergError};use async_trait::async_trait;
struct MinLengthValidator { min_length: usize,}
impl Plugin for MinLengthValidator { fn name(&self) -> &str { "min-length-validator" } fn version(&self) -> String { "1.0.0".to_string() } fn initialize(&self) -> Result<()> { Ok(()) } fn shutdown(&self) -> Result<()> { Ok(()) }}
#[async_trait]impl Validator for MinLengthValidator { async fn validate( &self, result: &ExtractedDocument, _config: &ExtractionConfig, ) -> Result<()> { if result.content.len() < self.min_length { return Err(XbergError::validation(format!( "Content too short: {} < {} characters", result.content.len(), self.min_length ))); } Ok(()) }
fn priority(&self) -> i32 { 100 }}from xberg import register_validator, ExtractedDocument, ExtractionConfig, ValidationError
class MinLengthValidator: def __init__(self, min_length: int = 100): self.min_length: int = min_length
def name(self) -> str: return "min_length_validator"
def version(self) -> str: return "1.0.0"
def priority(self) -> int: return 100
def validate(self, result: ExtractedDocument, config: ExtractionConfig) -> None: content_len: int = len(result.content) if content_len < self.min_length: raise ValidationError(f"Content too short: {content_len}")
def should_validate(self, result: ExtractedDocument, config: ExtractionConfig) -> bool: return True
def initialize(self) -> None: pass
def shutdown(self) -> None: pass
validator: MinLengthValidator = MinLengthValidator(min_length=100)register_validator(validator)import io.xberg.Xberg;import io.xberg.ExtractInputKind;import io.xberg.ExtractionResult;import io.xberg.ExtractedDocument;import io.xberg.ExtractInput;import io.xberg.ExtractionConfig;import io.xberg.IValidator;import io.xberg.ValidationException;import io.xberg.XbergRsException;
public class MinLengthValidatorExample implements IValidator { private final int minLength;
public MinLengthValidatorExample(int minLength) { this.minLength = minLength; }
@Override public String name() { return "min-length"; }
@Override public String version() { return "1.0.0"; }
@Override public void validate(ExtractedDocument result, ExtractionConfig config) throws Exception { if (result.content().length() < minLength) { throw new ValidationException( "Content too short: " + result.content().length() + " < " + minLength ); } }
@Override public boolean should_validate(ExtractedDocument _result, ExtractionConfig _config) throws Exception { return true; }
@Override public int priority() throws Exception { return 100; }
public static void main(String[] args) { try { Xberg.registerValidator(new MinLengthValidatorExample(100)); ExtractionResult output = Xberg.extract( ExtractInput.builder().withKind(ExtractInputKind.URI).withUri("document.pdf").build(), ExtractionConfig.builder().build() ); ExtractedDocument result = output.results().get(0); System.out.println("Validation passed!"); } catch (XbergRsException e) { // A ValidationException thrown from validate() is reported here, // wrapped by the native bridge, once the plugin call crosses back into Java. System.err.println("Validation failed: " + e.getMessage()); } }}using Xberg;
var validator = new MinLengthValidator(minLength: 100);ValidatorRegistry.RegisterValidator(validator);
public class MinLengthValidator : IValidator{ private readonly int _minLength;
public MinLengthValidator(int minLength = 100) { _minLength = minLength; }
public string Name => "min_length_validator"; public string Version => "1.0.0"; public int Priority => 100;
public void Validate(ExtractedDocument result, ExtractionConfig config) { var contentLength = result.Content.Length; if (contentLength < _minLength) throw new ValidationException($"Content too short: {contentLength}"); }
public bool ShouldValidate(ExtractedDocument result, ExtractionConfig config) => true; public void Initialize() { } public void Shutdown() { }}Quality Score Validator
Section titled “Quality Score Validator”use xberg::plugins::{Plugin, Validator};use xberg::{ExtractedDocument, ExtractionConfig, Result, XbergError};use async_trait::async_trait;
struct QualityValidator;
impl Plugin for QualityValidator { fn name(&self) -> &str { "quality-validator" } fn version(&self) -> String { "1.0.0".to_string() } fn initialize(&self) -> Result<()> { Ok(()) } fn shutdown(&self) -> Result<()> { Ok(()) }}
#[async_trait]impl Validator for QualityValidator { async fn validate( &self, result: &ExtractedDocument, _config: &ExtractionConfig, ) -> Result<()> { let score = result.metadata .additional .get("quality_score") .and_then(|v| v.as_f64()) .unwrap_or(0.0);
if score < 0.5 { return Err(XbergError::validation(format!( "Quality score too low: {:.2} < 0.50", score ))); }
Ok(()) }}from xberg import ExtractedDocument, ExtractionConfig, ValidationError, register_validator
class QualityValidator: def name(self) -> str: return "quality-validator"
def version(self) -> str: return "1.0.0"
def validate(self, result: ExtractedDocument, config: ExtractionConfig) -> None: score: float = result.quality_score or 0.0 if score < 0.5: raise ValidationError( f"Quality score too low: {score:.2f}" )
validator: QualityValidator = QualityValidator()register_validator(validator)import io.xberg.ExtractedDocument;import io.xberg.ExtractionConfig;import io.xberg.IValidator;import io.xberg.ValidationException;
IValidator qualityValidator = new IValidator() { @Override public String name() { return "quality-score"; }
@Override public String version() { return "1.0.0"; }
@Override public void validate(ExtractedDocument result, ExtractionConfig config) throws Exception { double score = result.qualityScore() != null ? result.qualityScore() : 0.0;
if (score < 0.5) { throw new ValidationException( String.format("Quality score too low: %.2f < 0.50", score) ); } }
@Override public boolean should_validate(ExtractedDocument _result, ExtractionConfig _config) throws Exception { return true; }
@Override public int priority() throws Exception { return 50; }};using Xberg;
public class QualityValidator : IValidator{ public string Name => "quality-validator"; public string Version => "1.0.0"; public int Priority => 100;
public void Validate(ExtractedDocument result, ExtractionConfig config) { var score = result.QualityScore ?? 0.0;
if (score < 0.5) throw new ValidationException($"Quality score too low: {score:F2}"); }
public bool ShouldValidate(ExtractedDocument result, ExtractionConfig config) => result.QualityScore.HasValue; public void Initialize() { } public void Shutdown() { }}
class Program{ static void Main() { var validator = new QualityValidator(); ValidatorRegistry.RegisterValidator(validator); }}Plugin Management
Section titled “Plugin Management”Listing
Section titled “Listing”Registry operations are type-specific. This example lists post-processors; the OCR, validator, embedding, reranker, tokenizer, and renderer registries expose the same operation.
List post-processors
from xberg import list_post_processors
def main() -> None: result = list_post_processors() print(result)
main()List post-processors
import { listPostProcessors } from "@xberg-io/xberg";function main() { const result = listPostProcessors(); console.log(result);}
void main();List post-processors
import { listPostProcessors } from "@xberg-io/xberg-wasm";function main() { const result = listPostProcessors(); console.log(result);}
void main();List post-processors
use xberg::list_post_processors;
fn main() { let result = list_post_processors(); println!("{:?}", result);}List post-processors
package main
import ( "fmt" xberg "github.com/xberg-io/xberg/packages/go")
func main() { result, err := xberg.ListPostProcessors() if err != nil { panic(err) } fmt.Printf("%+v\n", result)}List post-processors
import io.xberg.*;
public final class Example { public static void main(String[] args) throws Exception { var result = Xberg.listPostProcessors(); System.out.println(result); }}List post-processors
import io.xberg.*
fun main() { val result = Xberg.listPostProcessors() println(result)}List post-processors
using System;using Xberg;
var result = XbergConverter.ListPostProcessors();Console.WriteLine(result);List post-processors
import Xberg
let result = try Xberg.listPostProcessors()print(result)List post-processors
require "xberg"result = Xberg.list_post_processors()puts result.inspectList post-processors
<?php
declare(strict_types=1);
require_once __DIR__ . '/vendor/autoload.php';
use Xberg\Xberg;$result = Xberg::listPostProcessors();var_dump($result);List post-processors
result = Xberg.list_post_processors()IO.inspect(result)List post-processors
import 'dart:io';import 'package:xberg/xberg.dart';import 'package:xberg/src/xberg_bridge_generated/frb_generated.dart' show RustLib;Future<void> main() async { await RustLib.init(); try { final result = await XbergBridge.listPostProcessors(); stdout.writeln(result); } finally { RustLib.dispose(); }}List post-processors
const std = @import("std");const xberg = @import("xberg");
pub fn main() !void { const result = try xberg.list_post_processors(); std.debug.print("{any}\n", .{result});
}List post-processors
#include <assert.h>#include <stdint.h>#include <stdio.h>#include <stdlib.h>#include <string.h>#include "xberg.h"
int main(void) { char* result = xberg_list_post_processors(); (void)result; xberg_free_string(result); return EXIT_SUCCESS;}Unregistering
Section titled “Unregistering”Unregister from the registry that owns the plugin. The example registers a post-processor first so the removal has an observable target.
unregister_post_processor
use xberg::unregister_post_processor;
fn main() { let name = r#"test-processor"#; let _ = unregister_post_processor(name);}Clearing All
Section titled “Clearing All”Clear only the registry you intend to reset:
Clear all post-processors and verify list is empty
from xberg import clear_post_processors
def main() -> None: clear_post_processors()
main()Clear all post-processors and verify list is empty
import { clearPostProcessors } from "@xberg-io/xberg";function main() { const result = clearPostProcessors();}
void main();Clear all post-processors and verify list is empty
import { clearPostProcessors } from "@xberg-io/xberg-wasm";function main() { const result = clearPostProcessors();}
void main();Clear all post-processors and verify list is empty
use xberg::clear_post_processors;
fn main() { let _ = clear_post_processors();}Clear all post-processors and verify list is empty
package main
import ( xberg "github.com/xberg-io/xberg/packages/go")
func main() { err := xberg.ClearPostProcessors() if err != nil { panic(err) }}Clear all post-processors and verify list is empty
import io.xberg.*;
public final class Example { public static void main(String[] args) throws Exception { Xberg.clearPostProcessors(); }}Clear all post-processors and verify list is empty
import io.xberg.*
fun main() { PostProcessorBridge.clearAll()}Clear all post-processors and verify list is empty
using Xberg;
XbergConverter.ClearPostProcessors();Clear all post-processors and verify list is empty
import Xberg
try Xberg.clearPostProcessors()Clear all post-processors and verify list is empty
require "xberg"Xberg.clear_post_processors()Clear all post-processors and verify list is empty
<?php
declare(strict_types=1);
require_once __DIR__ . '/vendor/autoload.php';
use Xberg\Xberg;Xberg::clearPostProcessors();Clear all post-processors and verify list is empty
Xberg.clear_post_processors()Clear all post-processors and verify list is empty
import 'package:xberg/xberg.dart';import 'package:xberg/src/xberg_bridge_generated/frb_generated.dart' show RustLib;Future<void> main() async { await RustLib.init(); try { await XbergBridge.clearPostProcessors(); } finally { RustLib.dispose(); }}Clear all post-processors and verify list is empty
const std = @import("std");const xberg = @import("xberg");
pub fn main() !void { _ = try xberg.clear_post_processors();}Clear all post-processors and verify list is empty
#include <assert.h>#include <stdint.h>#include <stdio.h>#include <stdlib.h>#include <string.h>#include "xberg.h"
int main(void) { int32_t result = xberg_clear_post_processor(NULL); assert(result == 0 && "expected call to succeed"); return EXIT_SUCCESS;}Thread Safety
Section titled “Thread Safety”use std::collections::HashMap;use std::sync::{Arc, Mutex};use std::sync::atomic::{AtomicUsize, Ordering};use xberg::plugins::{Plugin, PostProcessor, ProcessingStage};use xberg::{ExtractedDocument, ExtractionConfig, Result, XbergError};use async_trait::async_trait;
struct StatefulPlugin { call_count: AtomicUsize, cache: Mutex<HashMap<String, String>>,}
impl Plugin for StatefulPlugin { fn name(&self) -> &str { "stateful-plugin" } fn version(&self) -> String { "1.0.0".to_string() }
fn initialize(&self) -> Result<()> { self.call_count.store(0, Ordering::Release); Ok(()) }
fn shutdown(&self) -> Result<()> { let count = self.call_count.load(Ordering::Acquire); println!("Plugin called {} times", count); Ok(()) }}
#[async_trait]impl PostProcessor for StatefulPlugin { async fn process( &self, result: &mut ExtractedDocument, _config: &ExtractionConfig ) -> Result<()> { self.call_count.fetch_add(1, Ordering::AcqRel);
let mut cache = self.cache.lock() .map_err(|_| XbergError::LockPoisoned("stateful-plugin cache".to_string()))?; cache.insert("last_mime".to_string(), result.mime_type.to_string());
Ok(()) }
fn processing_stage(&self) -> ProcessingStage { ProcessingStage::Middle }}import threadingfrom xberg import ExtractedDocument, ExtractionConfig
class StatefulPlugin: def __init__(self): self.lock: threading.Lock = threading.Lock() self.call_count: int = 0 self.cache: dict = {}
def name(self) -> str: return "stateful-plugin"
def version(self) -> str: return "1.0.0"
def processing_stage(self) -> str: return "early"
def process(self, result: ExtractedDocument, config: ExtractionConfig) -> None: with self.lock: self.call_count += 1 self.cache["last_mime"] = result.mime_type
def initialize(self) -> None: pass
def shutdown(self) -> None: passimport io.xberg.ExtractedDocument;import io.xberg.ExtractionConfig;import io.xberg.IPostProcessor;import java.util.concurrent.ConcurrentHashMap;import java.util.concurrent.atomic.AtomicInteger;
class StatefulPlugin implements IPostProcessor { // Use atomic types for simple counters private final AtomicInteger callCount = new AtomicInteger(0);
// Use concurrent collections for complex state private final ConcurrentHashMap<String, String> cache = new ConcurrentHashMap<>();
@Override public String name() { return "stateful-plugin"; }
@Override public String version() { return "1.0.0"; }
@Override public void process(ExtractedDocument result, ExtractionConfig config) { // Increment counter atomically callCount.incrementAndGet();
// Update cache (thread-safe) cache.put("last_mime", result.mimeType()); }
@Override public String processing_stage() { return "stateful-plugin"; }
@Override public boolean should_process(ExtractedDocument result, ExtractionConfig config) { return true; }
@Override public long estimated_duration_ms(ExtractedDocument result) { return 0; }
@Override public int priority() { return 50; }
public int getCallCount() { return callCount.get(); }}using Xberg;using System;using System.Collections.Concurrent;using System.Text.Json;
var processor = new StatefulPostProcessor();PostProcessorRegistry.RegisterPostProcessor(processor);Console.WriteLine("Post-processor registered");
public class StatefulPostProcessor : IPostProcessor{ private readonly object _lock = new(); private int _callCount = 0; private readonly ConcurrentDictionary<string, string> _cache = new();
public string Name => "stateful-plugin"; public string Version => "1.0.0"; public int Priority => 50; public ProcessingStage ProcessingStage => ProcessingStage.Middle;
public void Initialize() { } public void Shutdown() { }
public bool ShouldProcess(ExtractedDocument result, ExtractionConfig config) => true; public ulong EstimatedDurationMs(ExtractedDocument result) => 5;
public void Process(ExtractedDocument result, ExtractionConfig config) { lock (_lock) { _callCount++; _cache["last_mime"] = result.MimeType; } result.Metadata.Additional["call_count"] = JsonSerializer.SerializeToElement(_callCount); }}Best Practices
Section titled “Best Practices”Naming: Use kebab-case (my-custom-plugin), lowercase only, no spaces or special characters.
Logging
Section titled “Logging”import loggingfrom xberg import ExtractInput, ExtractionConfig
logger = logging.getLogger(__name__)
class MyPlugin: def name(self) -> str: return "my-plugin"
def version(self) -> str: return "1.0.0"
def supported_mime_types(self) -> list[str]: return ["application/x-custom"]
def initialize(self) -> None: logger.info(f"Initializing plugin: {self.name()}")
def shutdown(self) -> None: logger.info(f"Shutting down plugin: {self.name()}")
def extract(self, input: ExtractInput, config: ExtractionConfig) -> dict: logger.info(f"Extracting {input.mime_type} ({len(input.bytes or b'')} bytes)") result: dict = {"content": "", "mime_type": input.mime_type or "application/x-custom"} if not result["content"]: logger.warning("Extraction resulted in empty content") return resultuse xberg::plugins::{Plugin, DocumentExtractor};use xberg::{ExtractInput, ExtractedDocument, ExtractionConfig, Result};use async_trait::async_trait;use tracing::{debug, info, warn};
struct MyPlugin;
impl Plugin for MyPlugin { fn name(&self) -> &str { "my-plugin" }
fn version(&self) -> String { "1.0.0".to_string() }
fn initialize(&self) -> Result<()> { info!(plugin = self.name(), "initializing plugin"); Ok(()) }
fn shutdown(&self) -> Result<()> { info!(plugin = self.name(), "shutting down plugin"); Ok(()) }}
#[async_trait]impl DocumentExtractor for MyPlugin { #[tracing::instrument( name = "xberg::plugin_extract", level = "debug", skip_all, fields( plugin = self.name(), mime_type = tracing::field::Empty, input_len = tracing::field::Empty, ) )] async fn extract( &self, input: ExtractInput, _config: &ExtractionConfig, ) -> Result<ExtractedDocument> { let mime_type = input.mime_type.clone().unwrap_or_default(); let bytes = input.bytes.unwrap_or_default(); let span = tracing::Span::current(); span.record("mime_type", mime_type.as_str()); span.record("input_len", bytes.len()); debug!("extracting document");
let result = ExtractedDocument::default();
if result.content.is_empty() { warn!("extraction produced empty content"); }
Ok(result) }
fn supported_mime_types(&self) -> &[&str] { &["application/octet-stream"] }}import io.xberg.ExtractedDocument;import io.xberg.ExtractionConfig;import io.xberg.IPostProcessor;import java.util.logging.Logger;
class MyPlugin implements IPostProcessor { private static final Logger logger = Logger.getLogger(MyPlugin.class.getName());
@Override public String name() { return "my-plugin"; }
@Override public String version() { return "1.0.0"; }
@Override public void process(ExtractedDocument result, ExtractionConfig config) { logger.info("Processing " + result.mimeType() + " (" + result.content().length() + " bytes)");
// Processing...
if (result.content().isEmpty()) { logger.warning("Processing resulted in empty content"); } }
@Override public String processing_stage() { return "my-plugin"; }
@Override public boolean should_process(ExtractedDocument result, ExtractionConfig config) { return true; }
@Override public long estimated_duration_ms(ExtractedDocument result) { return 0; }
@Override public int priority() { return 50; }}using Xberg;using System;using System.Collections.Generic;using System.Text;
public class MyExtractorPlugin : IDocumentExtractor{ private readonly Action<string> _writeLog;
public MyExtractorPlugin(Action<string> writeLog) { _writeLog = writeLog ?? throw new ArgumentNullException(nameof(writeLog)); }
public string Name => "my-plugin"; public string Version => "1.0.0"; public int Priority => 50; public List<string> SupportedMimeTypes => new() { "text/plain" };
public void Initialize() { _writeLog($"INFO Initializing plugin: {Name}"); }
public void Shutdown() { _writeLog($"INFO Shutting down plugin: {Name}"); }
public bool CanHandle(string path, string mimeType) => mimeType == "text/plain";
public ExtractedDocument Extract(ExtractInput input, ExtractionConfig config) { _writeLog($"INFO Extracting {input.MimeType} ({input.Bytes?.Length ?? 0} bytes)"); var content = input.Bytes is null ? "" : Encoding.UTF8.GetString(input.Bytes); if (string.IsNullOrEmpty(content)) { _writeLog("WARN Extraction resulted in empty content"); } return new ExtractedDocument { Content = content, MimeType = input.MimeType ?? "text/plain", Metadata = new Metadata(), }; }}Testing
Section titled “Testing”from xberg import ExtractInput, ExtractionConfig
def test_custom_extractor() -> None: extractor = CustomJsonExtractor() json_data: bytes = b'{"message": "Hello, world!"}' input = ExtractInput(kind="bytes", bytes=json_data, mime_type="application/json") config = ExtractionConfig() result: dict = extractor.extract(input, config) assert "Hello, world!" in result["content"] assert result["mime_type"] == "application/json"#[cfg(test)]mod tests { use super::*;
#[tokio::test] async fn test_custom_extractor() { let extractor = CustomJsonExtractor;
let json_data = br#"{"message": "Hello, world!"}"#; let config = ExtractionConfig::default();
let result = extractor .extract(json_data, "application/json", &config) .await .expect("Extraction failed");
assert!(result.content.contains("Hello, world!")); assert_eq!(result.mime_type, "application/json"); }}import io.xberg.ExtractedDocument;import io.xberg.ExtractionConfig;import io.xberg.IPostProcessor;import io.xberg.Metadata;import org.junit.jupiter.api.Test;import static org.junit.jupiter.api.Assertions.*;
class PostProcessorTest { // process() returns void and ExtractedDocument is immutable, so a test // double captures the observed result on a field instead of returning a // mutated copy. static class WordCountProcessor implements IPostProcessor { long lastWordCount;
@Override public String name() { return "word-count"; }
@Override public String version() { return "1.0.0"; }
@Override public void process(ExtractedDocument result, ExtractionConfig config) { lastWordCount = result.content().split("\\s+").length; }
@Override public String processing_stage() { return "word-count"; }
@Override public boolean should_process(ExtractedDocument result, ExtractionConfig config) { return true; }
@Override public long estimated_duration_ms(ExtractedDocument result) { return 0; }
@Override public int priority() { return 50; } }
@Test void testWordCountProcessor() throws Exception { WordCountProcessor processor = new WordCountProcessor();
ExtractedDocument input = ExtractedDocument.builder() .withContent("Hello world test") .withMimeType("text/plain") .withMetadata(Metadata.builder().build()) .build();
processor.process(input, ExtractionConfig.builder().build());
assertEquals(3, processor.lastWordCount); }}using Xberg;using System;using System.Collections.Generic;using System.Text;
CustomExtractorTests.VerifyExtractsJsonContent();
public static class CustomExtractorTests{ public static void VerifyExtractsJsonContent() { var extractor = new CustomJsonExtractor(); var input = new ExtractInput { Kind = ExtractInputKind.Bytes, Bytes = Encoding.UTF8.GetBytes("{\"message\": \"Hello, world!\"}"), MimeType = "application/json", };
var result = extractor.Extract(input, new ExtractionConfig());
if (!result.Content.Contains("Hello, world!", StringComparison.Ordinal)) { throw new InvalidOperationException("Expected extracted JSON content was missing."); } if (result.MimeType != "application/json") { throw new InvalidOperationException($"Expected application/json, got {result.MimeType}."); } }}
public sealed class CustomJsonExtractor : IDocumentExtractor{ public string Name => "custom-json"; public string Version => "1.0.0"; public int Priority => 50; public List<string> SupportedMimeTypes => new() { "application/json" };
public void Initialize() { } public void Shutdown() { }
public bool CanHandle(string path, string mimeType) => mimeType == "application/json";
public ExtractedDocument Extract(ExtractInput input, ExtractionConfig config) { var content = input.Bytes is null ? "" : Encoding.UTF8.GetString(input.Bytes); return new ExtractedDocument { Content = content, MimeType = input.MimeType ?? "application/json", Metadata = new Metadata(), }; }}Complete Example: PDF Metadata Extractor
Section titled “Complete Example: PDF Metadata Extractor”from xberg import register_post_processor, ExtractedDocument, ExtractionConfigimport logging
logger = logging.getLogger(__name__)
class PdfMetadataExtractor: def __init__(self): self.processed_count: int = 0
def name(self) -> str: return "pdf_metadata_extractor"
def version(self) -> str: return "1.0.0"
def description(self) -> str: return "Logs PDF processing activity"
def processing_stage(self) -> str: return "early"
def should_process(self, result: ExtractedDocument, config: ExtractionConfig) -> bool: return result.mime_type == "application/pdf"
def process(self, result: ExtractedDocument, config: ExtractionConfig) -> None: self.processed_count += 1 logger.info(f"Processed PDF #{self.processed_count}")
def initialize(self) -> None: logger.info("PDF metadata extractor initialized")
def shutdown(self) -> None: logger.info(f"Processed {self.processed_count} PDFs")
processor: PdfMetadataExtractor = PdfMetadataExtractor()register_post_processor(processor)package main
import ( "log" "sync/atomic"
"github.com/xberg-io/xberg/packages/go")
type pdfMetadataObserver struct { processedCount atomic.Int64}
func (processor *pdfMetadataObserver) Name() string { return "pdf_metadata_observer" }func (processor *pdfMetadataObserver) Version() string { return "1.0.0" }func (processor *pdfMetadataObserver) Initialize() error { return nil }func (processor *pdfMetadataObserver) Shutdown() error { return nil }func (processor *pdfMetadataObserver) Priority() int32 { return 80 }func (processor *pdfMetadataObserver) ProcessingStage() xberg.ProcessingStage { return xberg.ProcessingStageEarly}func (processor *pdfMetadataObserver) ShouldProcess( result xberg.ExtractedDocument, _ xberg.ExtractionConfig,) bool { return result.MimeType == "application/pdf"}func (processor *pdfMetadataObserver) EstimatedDurationMs(_ xberg.ExtractedDocument) uint64 { return 1}func (processor *pdfMetadataObserver) Process( result xberg.ExtractedDocument, _ xberg.ExtractionConfig,) error { if result.Metadata != nil && result.Metadata.Title != nil { log.Printf("PDF title: %s", *result.Metadata.Title) } log.Printf("PDF content length: %d", len(result.Content)) processor.processedCount.Add(1) return nil}
func main() { processor := &pdfMetadataObserver{} if err := xberg.RegisterPostProcessor(processor); err != nil { log.Fatalf("register post-processor: %v", err) } defer func() { if err := xberg.UnregisterPostProcessor(processor.Name()); err != nil { log.Printf("unregister post-processor: %v", err) } log.Printf("PDFs processed: %d", processor.processedCount.Load()) }()
input := xberg.ExtractInputFromURI("document.pdf") result, err := xberg.Extract(*input, xberg.ExtractionConfig{}) if err != nil { log.Fatalf("extract PDF: %v", err) } log.Printf("PDF MIME type: %s", result.Results[0].MimeType)}import io.xberg.Xberg;import io.xberg.ExtractInputKind;import io.xberg.ExtractionResult;import io.xberg.ExtractedDocument;import io.xberg.ExtractInput;import io.xberg.ExtractionConfig;import io.xberg.IPostProcessor;import io.xberg.XbergRsException;import java.util.concurrent.atomic.AtomicInteger;import java.util.logging.Logger;
public class PdfMetadataExtractorExample implements IPostProcessor { private static final Logger logger = Logger.getLogger( PdfMetadataExtractorExample.class.getName() ); private final AtomicInteger processedCount = new AtomicInteger(0);
@Override public String name() { return "pdf-metadata-extractor"; }
@Override public String version() { return "1.0.0"; }
// Post-processors observe the extracted document; ExtractedDocument and // Metadata are immutable records, so this hook cannot inject new metadata // fields the way older PostProcessor implementations could. @Override public void process(ExtractedDocument result, ExtractionConfig config) throws Exception { if (!result.mimeType().equals("application/pdf")) { return; } processedCount.incrementAndGet(); logger.info("Processed PDF: " + processedCount.get()); }
@Override public String processing_stage() throws Exception { return "pdf-metadata-extractor"; }
@Override public boolean should_process(ExtractedDocument _result, ExtractionConfig _config) throws Exception { return true; }
@Override public long estimated_duration_ms(ExtractedDocument _result) throws Exception { return 0; }
@Override public int priority() throws Exception { return 50; }
public static void main(String[] args) { PdfMetadataExtractorExample pdfMetadata = new PdfMetadataExtractorExample(); try { Xberg.registerPostProcessor(pdfMetadata); logger.info("PDF metadata extractor initialized"); ExtractionResult output = Xberg.extract( ExtractInput.builder().withKind(ExtractInputKind.URI).withUri("document.pdf").build(), ExtractionConfig.builder().build() ); ExtractedDocument result = output.results().get(0); logger.info("Processed " + pdfMetadata.processedCount.get() + " PDFs"); } catch (XbergRsException e) { e.printStackTrace(); } }}using Xberg;using System;
var processor = new PdfMetadataExtractor();PostProcessorRegistry.RegisterPostProcessor(processor);
public class PdfMetadataExtractor : IPostProcessor{ private int _processedCount = 0;
public string Name => "pdf_metadata_extractor"; public string Version => "1.0.0"; public int Priority => 50; public ProcessingStage ProcessingStage => ProcessingStage.Early;
public bool ShouldProcess(ExtractedDocument result, ExtractionConfig config) => result.MimeType == "application/pdf";
public ulong EstimatedDurationMs(ExtractedDocument result) => 1;
public void Process(ExtractedDocument result, ExtractionConfig config) { _processedCount++; }
public void Initialize() { Console.WriteLine("PDF metadata extractor initialized"); }
public void Shutdown() { Console.WriteLine($"Processed {_processedCount} PDFs"); }}require 'xberg'
class PdfMetadataExtractor def initialize @count = 0 end
def call(result) return result unless result['mime_type'] == 'application/pdf' @count += 1 result['metadata'] ||= {} result['metadata']['pdf_order'] = @count result endend
extractor = PdfMetadataExtractor.newXberg.register_post_processor('pdf_metadata', extractor)
config = Xberg::ExtractionConfig.new( postprocessor: { enabled: true })
input = Xberg::ExtractInput.new(uri: 'report.pdf')result = Xberg.extract(input, config)puts "Metadata: #{result.results.first.metadata.inspect}"