This repository was archived by the owner on Oct 8, 2026. It is now read-only.
Repository navigation
Expand file tree
/
Copy pathsdk.py
More file actions
142 lines (120 loc) · 4.76 KB
/
Copy pathsdk.py
File metadata and controls
142 lines (120 loc) · 4.76 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
# sdk.py
from typing import Dict, Optional, List
# Direct imports are used because the entry point script fixes the path.
from processor.v1 import processor_pb2
from opencdc.v1 import opencdc_pb2
from config.v1 import parameter_pb2
# Global instance of the user-defined processor
PROCESSOR: Optional['Processor'] = None
# --- Memory Management (Simulation for testing) ---
_WASM_MEMORY = bytearray(1024 * 128) # 128KB
_OFFSET = 0
def _malloc(size):
global _OFFSET
ptr = _OFFSET
_OFFSET += size
return ptr
def _get_memory_view():
return _WASM_MEMORY
# --- Wasm Exported Functions ---
def export(name):
def decorator(func):
setattr(func, '_export_name', name)
return func
return decorator
@export("conduit.processor.v1.malloc")
def malloc(size: int) -> int:
"""
Allocates a block of memory of a specific size for the host to use.
This is the actual exported function that Conduit will call.
"""
return _malloc(size)
def write_to_memory(data: bytes) -> int:
"""Writes bytes to memory, returns a packed ptr-size uint64."""
size = len(data)
ptr = malloc(size)
_get_memory_view()[ptr:ptr+size] = data
return (ptr << 32) | size
def read_from_memory(packed_ptr: int) -> bytes:
"""Reads bytes from a packed ptr-size uint64."""
ptr = packed_ptr >> 32
size = packed_ptr & 0xFFFFFFFF
return _get_memory_view()[ptr:ptr+size]
# --- Processor Interface ---
class Processor:
def specification(self) -> 'Specification':
raise NotImplementedError
def configure(self, config_map: Dict[str, str]):
pass
def open(self):
pass
def process(self, records: List[opencdc_pb2.Record]) -> List['ProcessedRecord']:
raise NotImplementedError
def teardown(self):
pass
# --- SDK Data Structures ---
class Specification:
def __init__(self, name: str, version: str, summary: str, description: str, author: str, params: Dict[str, parameter_pb2.Parameter]):
self.name, self.version, self.summary, self.description, self.author, self.params = name, version, summary, description, author, params
class ProcessedRecord:
@staticmethod
def Ack(record: opencdc_pb2.Record) -> 'ProcessedRecord': return ProcessedRecord(record=record)
@staticmethod
def Filter() -> 'ProcessedRecord': return ProcessedRecord(filter=True)
@staticmethod
def Error(err: Exception) -> 'ProcessedRecord': return ProcessedRecord(error=err)
def __init__(self, record=None, filter=False, error=None):
self.record, self.filter, self.error = record, filter, error
@export("conduit.processor.v1.specification")
def specification(packed_ptr: int) -> int:
spec_model = PROCESSOR.specification()
response = processor_pb2.Specify.Response(
name=spec_model.name,
version=spec_model.version,
summary=spec_model.summary,
description=spec_model.description,
author=spec_model.author,
parameters=spec_model.params
)
return write_to_memory(response.SerializeToString())
@export("conduit.processor.v1.configure")
def configure(packed_ptr: int) -> int:
req_bytes = read_from_memory(packed_ptr)
request = processor_pb2.Configure.Request()
request.ParseFromString(req_bytes)
PROCESSOR.configure(request.parameters)
response = processor_pb2.Configure.Response()
return write_to_memory(response.SerializeToString())
@export("conduit.processor.v1.open")
def open_processor(packed_ptr: int) -> int:
PROCESSOR.open()
response = processor_pb2.Open.Response()
return write_to_memory(response.SerializeToString())
@export("conduit.processor.v1.process")
def process(packed_ptr: int) -> int:
req_bytes = read_from_memory(packed_ptr)
request = processor_pb2.Process.Request()
request.ParseFromString(req_bytes)
results = PROCESSOR.process(request.records)
response = processor_pb2.Process.Response()
for res in results:
processed_record_proto = processor_pb2.Process.ProcessedRecord()
if res.error:
processed_record_proto.error_record.error.message = str(res.error)
elif res.filter:
processed_record_proto.filter_record.SetInParent()
else: # ACK
processed_record_proto.single_record.CopyFrom(res.record)
response.records.append(processed_record_proto)
return write_to_memory(response.SerializeToString())
@export("conduit.processor.v1.teardown")
def teardown(packed_ptr: int) -> int:
PROCESSOR.teardown()
response = processor_pb2.Teardown.Response()
return write_to_memory(response.SerializeToString())
# --- SDK Entrypoint ---
def run(processor: Processor):
global PROCESSOR
if not isinstance(processor, Processor):
raise TypeError("argument must be an instance of sdk.Processor")
PROCESSOR = processor