clusternode.java
来自「jsr170接口的java实现。是个apache的开源项目。」· Java 代码 · 共 1,278 行 · 第 1/3 页
JAVA
1,278 行
record.writeString(null); write(record, ntDefs, true); record.writeChar('\0'); record.update(); setRevision(record.getRevision()); succeeded = true; } catch (JournalException e) { String msg = "Unable to create log entry: " + e.getMessage(); log.error(msg); } catch (Throwable e) { String msg = "Unexpected error while creating log entry."; log.error(msg, e); } finally { if (!succeeded && record != null) { record.cancelUpdate(); } } } /** * {@inheritDoc} */ public void reregistered(NodeTypeDef ntDef) { if (status != STARTED) { log.info("not started: nodetype operation ignored."); return; } Record record = null; boolean succeeded = false; try { record = journal.getProducer(PRODUCER_ID).append(); record.writeString(null); write(record, ntDef); record.writeChar('\0'); record.update(); setRevision(record.getRevision()); succeeded = true; } catch (JournalException e) { String msg = "Unable to create log entry: " + e.getMessage(); log.error(msg); } catch (Throwable e) { String msg = "Unexpected error while creating log entry."; log.error(msg, e); } finally { if (!succeeded && record != null) { record.cancelUpdate(); } } } /** * {@inheritDoc} */ public void unregistered(Collection qnames) { if (status != STARTED) { log.info("not started: nodetype operation ignored."); return; } Record record = null; boolean succeeded = false; try { record = journal.getProducer(PRODUCER_ID).append(); record.writeString(null); write(record, qnames, false); record.writeChar('\0'); record.update(); setRevision(record.getRevision()); succeeded = true; } catch (JournalException e) { String msg = "Unable to create log entry: " + e.getMessage(); log.error(msg); } catch (Throwable e) { String msg = "Unexpected error while creating log entry."; log.error(msg, e); } finally { if (!succeeded && record != null) { record.cancelUpdate(); } } } /** * {@inheritDoc} */ public void setListener(NodeTypeEventListener listener) { nodeTypeListener = listener; } /** * Workspace update channel. */ class WorkspaceUpdateChannel implements UpdateEventChannel { /** * Attribute name used to store record. */ private static final String ATTRIBUTE_RECORD = "record"; /** * Workspace name. */ private final String workspace; /** * Create a new instance of this class. * * @param workspace workspace name */ public WorkspaceUpdateChannel(String workspace) { this.workspace = workspace; } /** * {@inheritDoc} */ public void updateCreated(Update update) { if (status != STARTED) { log.info("not started: update create ignored."); return; } try { Record record = journal.getProducer(PRODUCER_ID).append(); update.setAttribute(ATTRIBUTE_RECORD, record); } catch (JournalException e) { String msg = "Unable to create log entry."; log.error(msg, e); } catch (Throwable e) { String msg = "Unexpected error while creating log entry."; log.error(msg, e); } } /** * {@inheritDoc} */ public void updatePrepared(Update update) { if (status != STARTED) { log.info("not started: update prepare ignored."); return; } Record record = (Record) update.getAttribute(ATTRIBUTE_RECORD); if (record == null) { String msg = "No record created."; log.warn(msg); return; } EventStateCollection events = update.getEvents(); ChangeLog changes = update.getChanges(); boolean succeeded = false; try { record.writeString(workspace); write(record, changes, events); record.writeChar('\0'); succeeded = true; } catch (JournalException e) { String msg = "Unable to create log entry: " + e.getMessage(); log.error(msg); } catch (Throwable e) { String msg = "Unexpected error while preparing log entry."; log.error(msg, e); } finally { if (!succeeded && record != null) { record.cancelUpdate(); update.setAttribute(ATTRIBUTE_RECORD, null); } } } /** * {@inheritDoc} */ public void updateCommitted(Update update) { if (status != STARTED) { log.info("not started: update commit ignored."); return; } Record record = (Record) update.getAttribute(ATTRIBUTE_RECORD); if (record == null) { String msg = "No record prepared."; log.warn(msg); return; } try { record.update(); setRevision(record.getRevision()); log.info("Appended revision: " + record.getRevision()); } catch (JournalException e) { String msg = "Unable to commit log entry."; log.error(msg, e); } catch (Throwable e) { String msg = "Unexpected error while committing log entry."; log.error(msg, e); } finally { update.setAttribute(ATTRIBUTE_RECORD, null); } } /** * {@inheritDoc} */ public void updateCancelled(Update update) { if (status != STARTED) { log.info("not started: update cancel ignored."); return; } Record record = (Record) update.getAttribute(ATTRIBUTE_RECORD); if (record != null) { record.cancelUpdate(); update.setAttribute(ATTRIBUTE_RECORD, null); } } /** * {@inheritDoc} */ public void setListener(UpdateEventListener listener) { if (workspace == null) { versionUpdateListener = listener; } else { wspUpdateListeners.remove(workspace); if (listener != null) { wspUpdateListeners.put(workspace, listener); } } } } /** * Workspace lock channel. */ class WorkspaceLockChannel implements LockEventChannel { /** * Workspace name. */ private final String workspace; /** * Create a new instance of this class. * * @param workspace workspace name */ public WorkspaceLockChannel(String workspace) { this.workspace = workspace; } /** * {@inheritDoc} */ public ClusterOperation create(NodeId nodeId, boolean deep, String owner) { if (status != STARTED) { log.info("not started: lock operation ignored."); return null; } try { Record record = journal.getProducer(PRODUCER_ID).append(); return new LockOperation(ClusterNode.this, workspace, record, nodeId, deep, owner); } catch (JournalException e) { String msg = "Unable to create log entry: " + e.getMessage(); log.error(msg); return null; } catch (Throwable e) { String msg = "Unexpected error while creating log entry."; log.error(msg, e); return null; } } /** * {@inheritDoc} */ public ClusterOperation create(NodeId nodeId) { if (status != STARTED) { log.info("not started: unlock operation ignored."); return null; } try { Record record = journal.getProducer(PRODUCER_ID).append(); return new LockOperation(ClusterNode.this, workspace, record, nodeId); } catch (JournalException e) { String msg = "Unable to create log entry: " + e.getMessage(); log.error(msg); return null; } catch (Throwable e) { String msg = "Unexpected error while creating log entry."; log.error(msg, e); return null; } } /** * {@inheritDoc} */ public void setListener(LockEventListener listener) { wspLockListeners.remove(workspace); if (listener != null) { wspLockListeners.put(workspace, listener); } } } /** * Invoked when a record starts. * * @param workspace workspace, may be <code>null</code> */ private void start(String workspace) { this.workspace = workspace; changeLog = new ChangeLog(); events = new ArrayList(); } /** * Process an update operation. * * @param operation operation to process */ private void process(ItemOperation operation) { operation.apply(changeLog); } /** * Process an event. * * @param event event */ private void process(EventState event) { events.add(event); } /** * Process a lock operation. * * @param nodeId node id * @param isDeep flag indicating whether lock is deep * @param owner lock owner */ private void process(NodeId nodeId, boolean isDeep, String owner) { LockEventListener listener = (LockEventListener) wspLockListeners.get(workspace); if (listener == null) { try { clusterContext.lockEventsReady(workspace); } catch (RepositoryException e) { String msg = "Unable to make lock listener for workspace " + workspace + " online: " + e.getMessage(); log.warn(msg); } listener = (LockEventListener) wspLockListeners.get(workspace); if (listener == null) { String msg = "Lock channel unavailable for workspace: " + workspace; log.error(msg); return; } } try { listener.externalLock(nodeId, isDeep, owner); } catch (RepositoryException e) { String msg = "Unable to deliver lock event: " + e.getMessage(); log.error(msg); } } /** * Process an unlock operation. * * @param nodeId node id */ private void process(NodeId nodeId) { LockEventListener listener = (LockEventListener) wspLockListeners.get(workspace); if (listener == null) { try { clusterContext.lockEventsReady(workspace); } catch (RepositoryException e) { String msg = "Unable to make lock listener for workspace " + workspace + " online: " + e.getMessage(); log.warn(msg); } listener = (LockEventListener) wspLockListeners.get(workspace); if (listener == null) { String msg = "Lock channel unavailable for workspace: " + workspace; log.error(msg); return; } } try { listener.externalUnlock(nodeId); } catch (RepositoryException e) { String msg = "Unable to deliver lock event: " + e.getMessage(); log.error(msg); } } /** * Process a namespace operation. * * @param oldPrefix old prefix. if <code>null</code> this is a fresh mapping * @param newPrefix new prefix. if <code>null</code> this is an unmap operation * @param uri uri to map prefix to */ private void process(String oldPrefix, String newPrefix, String uri) { if (namespaceListener == null) { String msg = "Namespace listener unavailable."; log.error(msg); return; } try { namespaceListener.externalRemap(oldPrefix, newPrefix, uri); } catch (RepositoryException e) { String msg = "Unable to deliver namespace operation: " + e.getMessage(); log.error(msg); } } /** * Process one or more node type registrations. * * @param c collection of node type definitions, if this is a register * operation; collection of <code>QName</code>s if this is * an unregister operation
⌨️ 快捷键说明
复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?